use super::*;
use saddle_runtime::profusegw::*;
use saddle_runtime::request_task::reserved_set::ReservedCollectionJoin;
use std::{
alloc::Layout,
pin::pin,
task::{Context, Poll, Waker},
};
pub(in crate::process) fn budget() -> saddle_admission::DeploymentResourceBudget {
budget_with_timeout(5000)
}
fn budget_with_timeout(request_timeout_ms: u64) -> saddle_admission::DeploymentResourceBudget {
let pending = saddle_admission::freeze_deployment_resource_budget(
1, 32, request_timeout_ms, 1, 1, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
)
.unwrap();
let (app, listener) = saddle_core::BootstrapRendezvousIssuer::issue()
.freeze_application(saddle_core::GeneratedApplicationFreezeSource::new(
"app",
b"descriptor",
&["route"],
))
.unwrap();
let listener = listener
.freeze_listener(saddle_core::ListenerStartupFreezeSource::new(
"app",
"127.0.0.1:8000".parse().unwrap(),
"127.0.0.1:9000".parse().unwrap(),
Duration::from_millis(request_timeout_ms),
))
.ok()
.unwrap();
let (whole, receipt) = saddle_core::pair_bootstrap_rendezvous(app, listener)
.ok()
.unwrap();
saddle_admission::bind_deployment_resource_budget_bootstrap(pending, whole, receipt)
.ok()
.unwrap()
}
fn dispatch(
_: saddle_boundary::ingress::AcceptedIngress,
_: (),
_: BusinessConfig<()>,
_: crate::database_capability::DatabaseRequest,
) -> std::future::Ready<Result<Vec<u8>>> {
std::future::ready(Ok(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}".to_vec()))
}
type Dispatch = fn(
saddle_boundary::ingress::AcceptedIngress,
(),
BusinessConfig<()>,
crate::database_capability::DatabaseRequest,
) -> std::future::Ready<Result<Vec<u8>>>;
struct HeldProcessCapacity(std::sync::Mutex<Option<
saddle_runtime::request_task::reserved_set::ReservedProcessTaskStorage>>);
impl saddle_core::ComponentLifecycle for HeldProcessCapacity {
fn name(&self) -> &'static str { "held-process-capacity" }
fn start(&self) -> saddle_core::LifecycleFuture<'_> {
Box::pin(async { Ok(()) })
}
fn shutdown(&self) -> saddle_core::LifecycleFuture<'_> {
Box::pin(async {
drop(self.0.lock().unwrap().take());
Ok(())
})
}
}
impl saddle_observability::root_diagnostic::RecordedComponentLifecycle for HeldProcessCapacity {
fn name(&self) -> &'static str { "held-process-capacity" }
fn start<'a>(&'a self, _: &'a saddle_core::ContextLabel,
_: &'a saddle_observability::SourceOutput)
-> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
Box::pin(async { Ok(()) })
}
fn shutdown<'a>(&'a self, application: &'a saddle_core::ContextLabel,
output: &'a saddle_observability::SourceOutput)
-> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
saddle_observability::root_diagnostic::RecordedComponentLifecycle::shutdown_with_primary(
self, application, output, None)
}
fn shutdown_with_primary<'a>(&'a self, _: &'a saddle_core::ContextLabel,
_: &'a saddle_observability::SourceOutput,
_: Option<saddle_core::DiagnosticOccurrence>)
-> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
Box::pin(async { drop(self.0.lock().unwrap().take()); Ok(()) })
}
}
struct AbortProcessWorker(Box<dyn Fn() + Send + Sync>);
impl saddle_core::ComponentLifecycle for AbortProcessWorker {
fn name(&self) -> &'static str { "abort-process-worker" }
fn start(&self) -> saddle_core::LifecycleFuture<'_> {
Box::pin(async {
(self.0)();
Err(process_error("saddle.process.test_forced_start_failure"))
})
}
fn shutdown(&self) -> saddle_core::LifecycleFuture<'_> {
Box::pin(async { Ok(()) })
}
}
impl saddle_observability::root_diagnostic::RecordedComponentLifecycle for AbortProcessWorker {
fn name(&self) -> &'static str { "abort-process-worker" }
fn start<'a>(&'a self, application: &'a saddle_core::ContextLabel,
output: &'a saddle_observability::SourceOutput)
-> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
Box::pin(async move {
(self.0)();
saddle_observability::root_diagnostic::RecordedSaddleError::component_result::<(), _>(
Err(std::io::Error::other("forced process start fault")), application, output,
saddle_observability::root_diagnostic::ComponentSourceKind::Start, None,
ErrorKind::Infrastructure, "saddle.process.test_forced_start_failure",
"forced process start fault")
})
}
fn shutdown<'a>(&'a self, application: &'a saddle_core::ContextLabel,
output: &'a saddle_observability::SourceOutput)
-> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
saddle_observability::root_diagnostic::RecordedComponentLifecycle::shutdown_with_primary(
self, application, output, None)
}
fn shutdown_with_primary<'a>(&'a self, _: &'a saddle_core::ContextLabel,
_: &'a saddle_observability::SourceOutput,
_: Option<saddle_core::DiagnosticOccurrence>)
-> saddle_observability::root_diagnostic::RecordedLifecycleFuture<'a> {
Box::pin(async { Ok(()) })
}
}
#[test]
fn process_entry_start_original_reaches_application_subprocess() {
const CHILD: &str = "SADDLE_PROCESS_ENTRY_START_ORIGINAL_CHILD";
if let Some(root) = std::env::var_os(CHILD) {
let root = std::path::PathBuf::from(root);
let shutdown_fault = std::env::var_os("SADDLE_PROCESS_ENTRY_SHUTDOWN_FAULT").is_some();
let output = saddle_observability::EmergencyDiagnostics::start_checked(
&saddle_observability::FileLoggingConfig::new(root.clone(),
saddle_observability::Rotation::Daily)).unwrap();
let selected = output.source_output().unwrap();
let handle = output.handle();
assert!(saddle_runtime::diagnostics::install_output(handle.clone()).is_ok());
let result = poll_ready(coordinate_profusegw_app_run(budget(), move |process| {
let (terminal, exit) = run_profusegw_owned_application_with_diagnostics(
process, output, move |lease| async move {
let startup = lease.take_database_startup_half().unwrap();
drop(startup);
let storage = if shutdown_fault { None }
else { Some(lease.try_reserved_task_storage().unwrap()) };
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let observer = saddle_observability::Observer::with_writer(
Default::default(), std::io::sink()).unwrap();
let mut app = saddle_runtime::Application::new();
app.install_lifecycle_observer_with_source_output(observer.clone(),
"app", selected);
app.register_recorded(observer.clone())?;
if let Some(storage) = storage {
app.register_recorded(HeldProcessCapacity(std::sync::Mutex::new(Some(storage))))?;
}
let dispatch: Dispatch = dispatch;
let entry = Arc::new(ProcessEntry::new(None, listener,
saddle_boundary::ingress::ProfuseGwListenerAdapter::new("app").unwrap(),
(), dispatch, BusinessConfig::unit(), None, lease, observer, None,
Some(handle)));
app.register_recorded_shared(entry.clone())?;
if shutdown_fault {
app.register_recorded(AbortProcessWorker(Box::new(move || {
entry.running.lock().unwrap().as_ref().unwrap().worker.abort();
})))?;
}
Ok(app)
});
assert_eq!(exit.shutdown, saddle_observability::DiagnosticShutdown::Finished);
std::future::ready(terminal)
}));
let error = result.unwrap_err();
let primary = match error {
saddle_runtime::profusegw::ProfuseGwCoordinatorFailure::Lifecycle(error) => error,
saddle_runtime::profusegw::ProfuseGwCoordinatorFailure::LifecycleWithCleanup { primary, .. } => primary,
_ => panic!("expected Application lifecycle failure"),
};
assert_eq!(primary.code(), if shutdown_fault {
"saddle.process.test_forced_start_failure"
} else { "saddle.process.request_storage_unavailable" });
assert!(!primary.source_unavailable());
assert!(!primary.cleanup_source_unavailable());
let id = primary.diagnostic().unwrap().id();
let rows: Vec<serde_json::Value> = std::fs::read_to_string(
root.join("saddle.emergency.log")).unwrap().lines()
.map(|line| serde_json::from_str(line).unwrap()).collect();
let contexts: Vec<_> = rows.iter().filter(|row| row["event"] == "request_error_original"
&& row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "context").collect();
assert_eq!(contexts.len(), 1, "Application must reuse the source receipt");
let context: serde_json::Value = serde_json::from_str(
contexts[0]["payload"].as_str().unwrap()).unwrap();
assert_eq!(context["stage"], "component_start");
assert_eq!(context["context"]["application"]["value"], "app");
let segments: Vec<_> = rows.iter().filter(|row|
row["event"] == "request_error_original"
&& row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "debug" && row["cause_depth"] == 0).collect();
assert_eq!(segments.last().unwrap()["state"], "field_end");
let original = segments.iter().map(|row| row["payload"].as_str().unwrap())
.collect::<String>();
if shutdown_fault {
let cleanup: Vec<_> = rows.iter().filter(|row| row["event"] == "request_error_original"
&& row["channel"] == "context"
&& row["occurrence"]["primary_diagnostic_id"] == id).collect();
assert_eq!(cleanup.len(), 1, "one linked JoinError source");
let cleanup_context: serde_json::Value = serde_json::from_str(
cleanup[0]["payload"].as_str().unwrap()).unwrap();
assert_eq!(cleanup_context["stage"], "component_cleanup");
let cleanup_id = cleanup[0]["occurrence"]["diagnostic_id"].as_u64().unwrap();
let joined = rows.iter().filter(|row| row["event"] == "request_error_original"
&& row["occurrence"]["diagnostic_id"] == cleanup_id
&& row["channel"] == "description" && row["cause_depth"] == 0)
.map(|row| row["payload"].as_str().unwrap()).collect::<String>();
assert!(joined.contains("cancelled"), "actual Tokio JoinError source");
assert!(rows.iter().any(|row| row["occurrence"]["diagnostic_id"] == cleanup_id
&& row["channel"] == "terminal" && row["state"] == "exposed_chain_complete"));
} else {
assert_eq!(original, "Admission(CapacityRejected)");
let cause = rows.iter().filter(|row| row["event"] == "request_error_original"
&& row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "debug" && row["cause_depth"] == 1)
.map(|row| row["payload"].as_str().unwrap()).collect::<String>();
assert_eq!(cause, format!("{:?}", saddle_admission::AdmissionError::CapacityRejected));
let description = rows.iter().filter(|row| row["event"] == "request_error_original"
&& row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "description" && row["cause_depth"] == 0)
.map(|row| row["payload"].as_str().unwrap()).collect::<String>();
assert_eq!(description,
format!("{}", saddle_admission::AdmissionError::CapacityRejected));
assert!(rows.iter().any(|row| row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "terminal" && row["state"] == "exposed_chain_complete"));
}
return;
}
for shutdown_fault in [false, true] {
let root = std::env::temp_dir().join(format!("saddle-process-entry-{}-{}",
std::process::id(), shutdown_fault));
std::fs::create_dir(&root).unwrap();
let mut command = std::process::Command::new(std::env::current_exe().unwrap());
command.arg("--exact")
.arg("process::reserved_entry::lifecycle_tests::process_entry_start_original_reaches_application_subprocess")
.env(CHILD, &root);
if shutdown_fault { command.env("SADDLE_PROCESS_ENTRY_SHUTDOWN_FAULT", "1"); }
assert!(command.status().unwrap().success());
std::fs::remove_file(root.join("saddle.emergency.log")).unwrap();
std::fs::remove_dir(root).unwrap();
}
}
#[test]
fn recorded_process_entry_sources_are_readable_before_return() {
const CHILD: &str = "SADDLE_RECORDED_PROCESS_ENTRY_MODE";
if let Ok(mode) = std::env::var(CHILD) {
let root = std::env::temp_dir().join(format!("saddle-recorded-process-{}-{}",
std::process::id(), mode));
std::fs::create_dir(&root).unwrap();
let mut output = saddle_observability::EmergencyDiagnostics::start_checked(
&saddle_observability::FileLoggingConfig::new(root.clone(),
saddle_observability::Rotation::Daily)).unwrap();
let selected = output.source_output().unwrap();
let handle = output.handle();
let target = output.target().to_owned();
let mode_for_factory = mode.clone();
let _ = poll_ready(coordinate_profusegw_app_run(budget(), move |process| {
let terminal = run_profusegw_owned_application(process, move |lease| async move {
drop(lease.take_database_startup_half().unwrap());
let held = if mode_for_factory.starts_with("start") {
Some(lease.try_reserved_task_storage().unwrap())
} else { None };
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let observer = saddle_observability::Observer::with_writer(
Default::default(), std::io::sink()).unwrap();
let application = saddle_core::ContextLabel::checked("process-recorded-test").unwrap();
let dispatch: Dispatch = dispatch;
let entry = ProcessEntry::new(None, listener,
saddle_boundary::ingress::ProfuseGwListenerAdapter::new("process-recorded-test").unwrap(),
(), dispatch, BusinessConfig::unit(), None, lease, observer, None,
Some(handle));
if mode_for_factory.starts_with("start") {
let failure = entry.recorded_start_source(&application, &selected).await.unwrap_err();
assert_eq!(failure.safe().code(), "saddle.process.request_storage_unavailable");
if mode_for_factory.ends_with("unconfirmed") {
assert!(failure.safe().source_unavailable());
let original = failure.original_if_unconfirmed().unwrap();
assert!(original.source().is_some(), "owned AdmissionError remains in chain");
assert!(std::fs::read_to_string(&target).unwrap().is_empty());
} else {
assert!(failure.original_if_unconfirmed().is_none());
let id = failure.safe().diagnostic().unwrap().id();
let rows: Vec<serde_json::Value> = std::fs::read_to_string(&target).unwrap()
.lines().map(|row| serde_json::from_str(row).unwrap()).collect();
let original = |channel| rows.iter().filter(|row|
row["event"] == "request_error_original"
&& row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == channel && row["cause_depth"] == 0)
.map(|row| row["payload"].as_str().unwrap()).collect::<String>();
assert_eq!(original("description"),
format!("{}", saddle_admission::AdmissionError::CapacityRejected));
assert_eq!(original("debug"),
format!("{:?}", saddle_admission::AdmissionError::CapacityRejected));
assert!(rows.iter().any(|row| row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "description" && row["cause_depth"] == 1
&& row["payload"].as_str().is_some_and(|text| text.contains("capacity"))));
let context: serde_json::Value = serde_json::from_str(&original("context")).unwrap();
assert_eq!(context["stage"], "component_start");
assert_eq!(context["context"]["application"]["value"], "process-recorded-test");
assert!(rows.iter().any(|row| row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "terminal" && row["state"] == "exposed_chain_complete"));
let recovered: std::result::Result<(), _> = if failure.safe().code()
== "saddle.process.request_storage_unavailable" { Ok(()) }
else { Err(failure) };
assert!(recovered.is_ok(), "confirmed original survives legal recovery");
}
} else {
assert!(entry.recorded_start_source(&application, &selected).await.is_ok());
if mode_for_factory.starts_with("cleanup") {
entry.running.lock().unwrap().as_ref().unwrap().worker.abort();
}
let primary = saddle_core::Diagnostic::capture(
saddle_core::DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::DiagnosticCause::new(saddle_core::DiagnosticStage::ShutdownComponent,
saddle_core::DiagnosticCode::new("process.test_primary").unwrap()),
).occurrence();
let result = entry.recorded_cleanup_source(&application, &selected,
Some(primary.clone())).await;
if mode_for_factory.starts_with("cleanup") {
let failure = result.unwrap_err();
assert_eq!(failure.safe().code(), "saddle.process.entry_failed");
if mode_for_factory.ends_with("unconfirmed") {
assert!(failure.safe().source_unavailable());
let original = failure.original_if_unconfirmed().unwrap();
assert!(original.source().is_some(), "owned JoinError remains in chain");
assert!(std::fs::read_to_string(&target).unwrap().is_empty());
} else {
assert!(failure.original_if_unconfirmed().is_none());
let id = failure.safe().diagnostic().unwrap().id();
assert_eq!(serde_json::to_value(failure.safe().diagnostic().unwrap()).unwrap()
["primary_diagnostic_id"], serde_json::to_value(primary).unwrap()["diagnostic_id"]);
let rows: Vec<serde_json::Value> = std::fs::read_to_string(&target).unwrap()
.lines().map(|row| serde_json::from_str(row).unwrap()).collect();
let original = |channel| rows.iter().filter(|row|
row["event"] == "request_error_original"
&& row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == channel && row["cause_depth"] == 0)
.map(|row| row["payload"].as_str().unwrap()).collect::<String>();
assert!(original("description").contains("cancelled"));
assert!(original("debug").contains("Cancelled"));
assert!(rows.iter().any(|row| row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "description" && row["cause_depth"] == 1
&& row["payload"].as_str().is_some_and(|text| text.contains("cancelled"))));
let context: serde_json::Value = serde_json::from_str(&original("context")).unwrap();
assert_eq!(context["stage"], "component_cleanup");
assert_eq!(context["context"]["application"]["value"], "process-recorded-test");
assert!(rows.iter().any(|row| row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "terminal" && row["state"] == "exposed_chain_complete"));
let recovered: std::result::Result<(), _> = if failure.safe().code()
== "saddle.process.entry_failed" { Ok(()) }
else { Err(failure) };
assert!(recovered.is_ok(), "cleanup original survives legal recovery");
}
} else {
assert!(result.is_ok());
assert!(std::fs::read_to_string(&target).unwrap().is_empty());
}
}
drop((entry, held));
Err(process_error("saddle.process.test_finished"))
});
std::future::ready(terminal)
}));
let deadline = std::time::Instant::now() + Duration::from_secs(5);
while output.shutdown() == saddle_observability::DiagnosticShutdown::Pending
&& std::time::Instant::now() < deadline { std::thread::yield_now(); }
assert_eq!(output.shutdown(), saddle_observability::DiagnosticShutdown::Finished);
drop(output);
std::fs::remove_dir_all(root).unwrap();
return;
}
for mode in ["start", "cleanup", "normal", "start-unconfirmed", "cleanup-unconfirmed"] {
let mut command = if mode.ends_with("unconfirmed") {
let mut command = std::process::Command::new("bash");
command.arg("-c").arg("trap '' XFSZ; ulimit -f 0; exec \"$@\"")
.arg("process-unconfirmed-child")
.arg(std::env::current_exe().unwrap());
command
} else { std::process::Command::new(std::env::current_exe().unwrap()) };
let output = command.arg("--exact")
.arg("process::reserved_entry::lifecycle_tests::recorded_process_entry_sources_are_readable_before_return")
.env(CHILD, mode).output().unwrap();
assert!(output.status.success(), "{mode}: {}", String::from_utf8_lossy(&output.stderr));
}
}
#[test]
fn formal_factory_matching_join_last_reference_and_partial_write() {
let mut coordinator = pin!(coordinate_profusegw_app_run(budget(), |process| {
std::future::ready(run_profusegw_owned_application(
process,
|lease| async move {
let startup = lease.take_database_startup_half().unwrap();
let observer = saddle_observability::Observer::with_writer(
Default::default(),
std::io::sink(),
)
.unwrap();
let application = saddle_core::ContextLabel::checked("app").unwrap();
let dispatch: Dispatch = dispatch;
let prepared = prepare::<(), (), _, _>("app", &lease, &dispatch).unwrap();
let Prepared {
mut normal,
mut rejected,
indirect,
..
} = prepared;
for mode in 0..5 {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let mut client = TcpStream::connect(listener.local_addr().unwrap())
.await
.unwrap();
let (socket, _) = listener.accept().await.unwrap();
let input = NormalInput {
factory: Some(consumer_storage::NormalFactory {
socket,
input_planner: None,
preparsed: None,
adapter: saddle_boundary::ingress::ProfuseGwListenerAdapter::new("app")
.unwrap(),
deployment: Some(()),
business: Some(BusinessConfig::unit()),
ingress_token: None,
dispatch: Arc::new(dispatch),
observer: observer.clone(),
admission: None,
database: None,
source_output: consumer_storage::FactorySource::MissingForTest,
}),
retained: EntryOwner::new(),
};
let deadline = ProfuseGwIngressDeadline::start(Duration::from_millis(if mode == 3 { 30 } else { 5000 })).unwrap();
let outcome = lease.try_reserved_dispatch_at(
application.clone(),
None,
body_layout::<(), (), Dispatch, std::future::Ready<Result<Vec<u8>>>>(),
&indirect,
input,
factory::<(), (), _, _>,
&deadline,
);
let ReservedDispatchOutcome::Ready {
root,
future,
ticket,
} = outcome
else {
panic!("actual factory must fit")
};
let held = root.view(saddle_core::RequestViewPhase::Reading).unwrap();
let abort = normal
.spawn(future, ticket)
.unwrap_or_else(|_| panic!("original task slot"));
if mode == 0 {
abort.abort();
} else if mode == 4 {
let body = br#"{"target":{"app":"app","interfaceId":"route"},"profuseGwContext":{"userInfo":{"userId":"u"},"traceInfo":{"rpcId":"0"},"ldcInfo":{"zone":"z","idc":"i","env":"test"}},"requestData":{}}"#;
let head = format!("POST /saddle/v1/ingress/profusegw/invoke HTTP/1.1\r\nContent-Type: application/json\r\nContent-Length: {}\r\nX-Request-Id: r\r\nX-Call-Id: c\r\n\r\n",body.len());
client.write_all(head.as_bytes()).await.unwrap();
tokio::time::sleep(Duration::from_millis(5)).await;
assert!(lease.resource_snapshot().unwrap().charged >= MAX_HEAD_BYTES);
client.write_all(body).await.unwrap();
} else {
client.write_all(b"POST / HTTP/1.1\r\n").await.unwrap();
for _ in 0..3 {
tokio::task::yield_now().await;
}
if mode == 1 {
abort.abort();
} else if mode == 2 {
client.shutdown().await.unwrap();
}
}
while !abort.is_finished() {
tokio::task::yield_now().await;
}
assert_eq!(normal.len(), 1);
if mode==4 {
let before=lease.resource_snapshot().unwrap().framework_charged;
let ProfuseGwCoordinatorAdmissionOutcome::Ready(next,_) = lease.try_admit()
else {panic!("physical DB return permits a successor before task collection")};
next.cancel();
assert!(lease.resource_snapshot().unwrap().framework_charged>=before);
assert_eq!(normal.len(),1);
} else {
assert!(matches!(lease.try_admit(),ProfuseGwCoordinatorAdmissionOutcome::CapacityRejected(_)));
}
let ReservedCollectionJoin::Matched(joined) = normal.join_next().await.unwrap()
else {
panic!("matching ticket")
};
finish_normal(joined, &observer, None).await;
if mode >= 3 {
let mut response = Vec::new();
tokio::time::timeout(Duration::from_secs(1),client.read_to_end(&mut response)).await.unwrap().unwrap();
if mode == 3 {
assert_eq!(response, b"HTTP/1.1 408 Error\r\nContent-Length: 0\r\nConnection: close\r\n\r\n");
} else {
assert_eq!(response, b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}");
assert_eq!(lease.resource_snapshot().unwrap().charged, 0);
}
}
drop(root);
let before_last_view=lease.resource_snapshot().unwrap().framework_charged;
let ProfuseGwCoordinatorAdmissionOutcome::Ready(next, _) = lease.try_admit()
else {panic!("returned DB may be reused while the original root remains charged")};
next.cancel();
assert_eq!(lease.resource_snapshot().unwrap().framework_charged,before_last_view);
if mode == 1 {
let (mut write, mut read) = tokio::io::duplex(4);
let mut retained = None;
let mut future = Box::pin(deliver(
&mut write,
b"abcdefgh",
tokio::time::Instant::now() + Duration::from_secs(60),
held.clone(),
None,
&mut retained,
));
assert!(matches!(
future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop())),
Poll::Pending
));
let mut bytes = [0; 4];
read.read_exact(&mut bytes).await.unwrap();
assert_eq!(&bytes, b"abcd");
drop(future);
let ReservedDeliveryOutcome::Failed { axes, failure } =
retained.take().unwrap()
else {
panic!("cancel retains failure")
};
assert_eq!(axes.bytes_written, Some(4));
assert!(matches!(
axes.operation,
saddle_core::OperationOutcome::Cancelled
));
let (_, source) = failure.into_parts();
source
.finish(
&held,
None,
RootOutcomeFacts {
axes,
..Default::default()
},
)
.ok()
.unwrap();
}
drop(held);
assert!(lease.resource_snapshot().unwrap().framework_charged<before_last_view,"physical last view owns its bytes until drop");
let ProfuseGwCoordinatorAdmissionOutcome::Ready(owner, _) = lease.try_admit()
else {
panic!("last view must refund")
};
owner.cancel();
println!(
"S_REAL_FACTORY mode={mode} no_refund_before_join no_refund_before_last_view refund_after_last_view PASS"
);
}
println!(
"S_REAL_COLLECTION capacity={} ticket_entry={} rejection_body={}",
normal.capacity(),
Layout::new::<(
tokio::task::Id,
saddle_runtime::request_task::reserved::ReservedTaskTicket<
Owner<
(),
(),
fn(
saddle_boundary::ingress::AcceptedIngress,
(),
BusinessConfig<()>,
crate::database_capability::DatabaseRequest,
)
-> std::future::Ready<Result<Vec<u8>>>,
>,
>
)>()
.size(),
rejected_layout().size()
);
// Actual rejection factory and body within this 16-client fixture.
// Only real byte capacity may reject, not this fixture count.
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let mut slots = Vec::new();
let mut peers = Vec::new();
for _ in 0..16 {
let peer = TcpStream::connect(listener.local_addr().unwrap())
.await
.unwrap();
let (socket, _) = listener.accept().await.unwrap();
let input = RejectedInput {
socket: Some(socket),
code: 503,
admission: None,
retained: EntryOwner::new(),
observer: observer.clone(),
output: None,
};
match lease.try_reserved_rejection(
application.clone(),
None,
rejected_layout(),
&[],
input,
rejected_factory,
) {
Ok(pair) => {
slots.push(pair);
peers.push(peer);
}
Err((input, _, error)) => {
let _facts = record_unrooted_failure(
&application,
None,
UnrootedEntryError::Rejection(&error),
);
drop(input);
break;
}
}
}
assert!(!slots.is_empty());
let fitted = slots.len();
// Grow the actual process table while every complete task is
// retained. The table's old backing remains charged during grow.
rejected.try_capacity(fitted).unwrap();
let held = slots
.last()
.unwrap()
.0
.view(saddle_core::RequestViewPhase::Reading)
.unwrap();
for (root, future, ticket) in slots {
let handle = rejected
.spawn(future, ticket)
.unwrap_or_else(|_| panic!("reserved rejection collection"));
handle.abort();
drop(root);
}
while let Some(join) = rejected.join_next().await {
let ReservedCollectionJoin::Matched(joined) = join else {
panic!("matching rejected ticket")
};
finish_rejected(joined).await;
}
assert_eq!(rejected.len(), 0);
drop((held, peers));
println!(
"S_REAL_REJECTION fitted={fitted} slot_limit=16 retained_capacity={} owner={} body={} matching_abort_cleanup PASS",
rejected.capacity(),
Layout::new::<RejectedInput>().size(),
rejected_layout().size()
);
drop((normal, rejected, startup, lease));
assert!(
std::process::Command::new("kill")
.args(["-TERM", &std::process::id().to_string()])
.status()
.unwrap()
.success()
);
Ok(saddle_runtime::Application::new())
},
))
}));
let Poll::Ready(Ok(())) = coordinator
.as_mut()
.poll(&mut Context::from_waker(Waker::noop()))
else {
panic!("original finalization must reach ZERO")
};
}
struct ListenerInput { id:Option<u64>, product_id:Option<u64>, label:Option<String>, rows:Option<Vec<Vec<String>>> }
impl saddle_admission::DecodeInput for ListenerInput {
fn inspect(value:saddle_admission::InputValue<'_>)->std::result::Result<(),saddle_admission::InputDecodeFailure> {
for name in ["id", "product_id", "label", "rows"] {
if let Some(field) = value.field(name) {
match name {
"id" | "product_id" => <Option<u64> as saddle_admission::DecodeInput>::inspect(field),
"label" => <Option<String> as saddle_admission::DecodeInput>::inspect(field),
_ => <Option<Vec<Vec<String>>> as saddle_admission::DecodeInput>::inspect(field),
}.map_err(|error| error.at(name))?;
}
}
Ok(())
}
fn plan(value:saddle_admission::InputValue<'_>,plan:&mut saddle_admission::InputPlan)->std::result::Result<(),saddle_admission::AdmissionError> {
plan.field::<Option<u64>>(value,"id")?;
plan.field::<Option<u64>>(value,"product_id")?;
plan.field::<Option<String>>(value,"label")?;
plan.field::<Option<Vec<Vec<String>>>>(value,"rows")
}
fn decode(value:saddle_admission::InputValue<'_>,builder:&mut saddle_admission::InputBuilder)->std::result::Result<Self,saddle_admission::AdmissionError> {
Ok(Self {id:builder.field(value,"id")?,product_id:builder.field(value,"product_id")?,label:builder.field(value,"label")?,rows:builder.field(value,"rows")?})
}
}
static INPUT_PLANNED: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false);
static LARGE_RESPONSE_TEXT: [u8; 16_384] = [b'R'; 16_384];
fn complete_original_with(rows: &str, marker: &str) -> bool {
let records: Vec<serde_json::Value> = rows.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "request_error_original")
.collect();
let Some(source) = records.iter().find(|row| row["payload"].as_str().is_some_and(|text| text.contains(marker))) else {
return false;
};
let occurrence = &source["occurrence"];
let matching: Vec<&serde_json::Value> = records.iter().filter(|row| &row["occurrence"] == occurrence).collect();
matching.iter().enumerate().all(|(sequence, row)| row["sequence"].as_u64() == Some(sequence as u64))
&& matching.last().is_some_and(|row| row["channel"] == "terminal" && row["state"] == "exposed_chain_complete")
}
fn complete_original_occurrences(rows: &str, marker: &str) -> usize {
let records: Vec<serde_json::Value> = rows.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "request_error_original")
.collect();
let mut occurrences = std::collections::HashSet::new();
for source in records.iter().filter(|row| row["payload"].as_str().is_some_and(|text| text.contains(marker))) {
let occurrence = &source["occurrence"];
assert!(!occurrence.is_null(), "source occurrence");
if !occurrences.insert(occurrence.to_string()) { continue; }
let matching: Vec<_> = records.iter().filter(|row| &row["occurrence"] == occurrence).collect();
assert!(matching.iter().enumerate().all(|(sequence, row)| row["sequence"].as_u64() == Some(sequence as u64)));
assert!(matching.last().is_some_and(|row| row["channel"] == "terminal" && row["state"] == "exposed_chain_complete"));
}
occurrences.len()
}
enum ResponseCapacityCode { Unused }
impl crate::ingress::ProfuseGwCode for ResponseCapacityCode {
const REGISTERED_CODES: &'static [&'static str] = &["UNUSED"];
fn stable_code(&self) -> &'static str { "UNUSED" }
}
fn plan_listener_input(
input: &saddle_boundary::ingress::AcceptedIngress,
recorder: &mut dyn saddle_admission::InputDecodeRecorder,
) -> std::result::Result<usize,saddle_admission::InputDecodeFailure> {
let peak = crate::programming::plan_accepted_profusegw::<ListenerInput>(input, recorder)?;
INPUT_PLANNED.store(true, Ordering::Release);
Ok(peak)
}
fn formal_listener_case(mode: u8) {
INPUT_PLANNED.store(false, Ordering::Release);
let mut coordinator = pin!(coordinate_profusegw_app_run(budget(), |process| {
std::future::ready(run_profusegw_owned_application(
process,
|lease| async move {
let startup = lease.take_database_startup_half().unwrap();
let observer = saddle_observability::Observer::with_writer(
Default::default(),
std::io::sink(),
)
.unwrap();
let mut diagnostic = if matches!(mode, 10 | 11 | 12 | 13 | 14 | 15 | 16 | 17 | 18 | 19 | 20 | 21 | 22 | 23 | 24 | 25 | 26 | 28 | 29 | 30 | 31 | 32 | 33 | 34 | 35 | 36 | 37 | 38 | 40 | 41) {
let directory = std::env::temp_dir().join(format!(
"saddle-i027-entry-facts-{}-{mode}", std::process::id()));
std::fs::create_dir_all(&directory).unwrap();
let writer = saddle_observability::EmergencyDiagnostics::start_checked(
&saddle_observability::FileLoggingConfig::new(
&directory, saddle_observability::Rotation::Daily)).unwrap();
Some((writer, directory))
} else { None };
let diagnostic_handle = diagnostic.as_ref().map(|(writer, _)| writer.handle());
let factory_source = diagnostic.as_ref()
.and_then(|(writer, _)| writer.source_output())
.map(consumer_storage::FactorySource::Selected)
.unwrap_or(consumer_storage::FactorySource::MissingForTest);
let diagnostic_path = diagnostic.as_ref().map(|(_, directory)| directory.join("saddle.emergency.log"));
let (database_entry, database_state) = if matches!(mode, 34 | 35 | 40 | 41) {
use saddle_runtime::startup_assembly::StartupDbPoolFactory as _;
let connection_file = std::path::PathBuf::from(
std::env::var("SADDLE_I_DB_CONNECTION_FILE").unwrap());
let connection: serde_json::Value = serde_json::from_slice(
&std::fs::read(&connection_file).unwrap()).unwrap();
let mappings = connection_file.parent().unwrap().join("empty-mappings");
std::fs::create_dir_all(&mappings).unwrap();
let config = saddle_db::DatabaseConfig::new(connection["url"].as_str().unwrap())
.max_connections(1).name_mapping_directory(mappings);
let factory = saddle_db::internal::StartupManagedDatabaseFactory::new(
Some(config), observer.clone());
let factory = if matches!(mode, 34 | 40 | 41) {
factory.with_diagnostic_output(diagnostic_handle.as_ref().unwrap().clone())
} else { factory };
let owner = factory
.construct(saddle_admission::DbCreditProfile {
connections: 1, operations: 1,
}).await.unwrap();
let (entry, state) = ManagedDatabaseEntry::new(owner, startup,
"app".to_owned(), diagnostic_handle.clone());
saddle_runtime::ComponentLifecycle::start(&entry).await.unwrap();
(Some(entry), Some(state))
} else {
drop(startup);
(None, None)
};
let entered = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let seen = entered.clone();
let started=Arc::new(tokio::sync::Notify::new());
let started_signal=started.clone();
let release=Arc::new(tokio::sync::Notify::new());
let signal=release.clone();
let order=Arc::new(std::sync::Mutex::new(Vec::new()));
let recorded=order.clone();
let source_barrier=Arc::new(std::sync::atomic::AtomicUsize::new(0));
let source_barrier_dispatch=source_barrier.clone();
let dispatch = move |input:saddle_boundary::ingress::AcceptedIngress, _:(), _:BusinessConfig<()>, request:crate::database_capability::DatabaseRequest| {
let diagnostic_path = diagnostic_path.clone();
let source_barrier=source_barrier_dispatch.clone();
seen.fetch_add(1,Ordering::SeqCst);
started_signal.notify_one();
let saddle_boundary::ingress::AcceptedRequestData::Managed(document)=input.request_data else { panic!("formal input must be metered") };
let typed=document.root().field("requestData").unwrap().read_only::<ListenerInput>().unwrap();
drop(document);
let id=typed.get().id.unwrap_or(0);
recorded.lock().unwrap().push(id);
let signal=signal.clone();
async move {
let mut request = request;
if mode==5 {signal.notified().await;}
if mode==13 && id==0 { signal.notified().await; }
if mode==12 && id==0 {
// Finite real computation on the existing Runtime;
// bounded slices let the listener keep sampling.
let until = std::time::Instant::now() + Duration::from_millis(3900);
while std::time::Instant::now() < until {
let slice = std::time::Instant::now() + Duration::from_millis(10);
while std::time::Instant::now() < slice { std::hint::spin_loop(); }
tokio::task::yield_now().await;
}
}
assert_eq!(typed.get().label.as_deref(),Some("held"));
assert_eq!(typed.get().rows.as_ref().unwrap()[0][0],"copy");
if mode==9 {
signal.notified().await;
let encoding=request.response_encoding_context();
drop(request);
let value=std::str::from_utf8(&LARGE_RESPONSE_TEXT).unwrap();
let response=crate::ingress::ProfuseGwResponse::<_,ResponseCapacityCode>::success(value);
let wire=encoding.encode(move |memory|crate::profusegw_http::encode_profusegw_http1_managed(&response,memory)).await?;
return Ok(wire.as_slice().to_vec());
}
if matches!(mode, 15 | 16 | 18) {
let source = if mode == 18 {
crate::programming::ProfuseGwDispatchError::IdentityMismatch
} else {
crate::programming::ProfuseGwDispatchError::InterfaceNotFound
};
let public = if mode != 16 {
let source = request.record_dispatch_at_source(source);
request.dispatch_recorder().record_result::<()>(Err(source)).unwrap_err()
} else {
let public = request.dispatch_recorder().record_result::<()>(Err(source)).unwrap_err();
assert_eq!(public.code(), "saddle.diagnostic.source_unavailable");
public
};
let rows = std::fs::read_to_string(diagnostic_path.as_ref().unwrap()).unwrap();
assert_eq!(complete_original_with(&rows, if mode == 18 { "IdentityMismatch" } else { "InterfaceNotFound" }), mode != 16,
"only source branches with a receipt may claim a written original");
drop(request);
return Err(public);
}
if mode == 17 {
let external = std::io::Error::new(std::io::ErrorKind::InvalidData, "external decode cause");
let category = request.record_external_dispatch_at_source(
&external, crate::programming::ProfuseGwDispatchError::RequestDataInvalid);
let rows = std::fs::read_to_string(diagnostic_path.as_ref().unwrap()).unwrap();
assert!(complete_original_with(&rows, "external decode cause"),
"external source must be written before classification");
let public = request.dispatch_recorder().record_result::<()>(Err(category)).unwrap_err();
drop(request);
return Err(public);
}
if mode == 28 {
let external = std::io::Error::new(std::io::ErrorKind::InvalidData,
format!("closed writer long external cause: {}", "原文".repeat(5000)));
let category = request.record_external_dispatch_at_source(
&external, crate::programming::ProfuseGwDispatchError::RequestDataInvalid);
let rows = std::fs::read_to_string(diagnostic_path.as_ref().unwrap()).unwrap();
assert!(complete_original_with(&rows, "closed writer long external cause"),
"every source segment and end marker must survive worker shutdown");
let public = request.dispatch_recorder().record_result::<()>(Err(category)).unwrap_err();
drop(request);
return Err(public);
}
if matches!(mode, 29 | 30) {
let category = crate::programming::ProfuseGwDispatchError::InterfaceNotFound;
request.record_dispatch_at_source(category);
request.record_dispatch_at_source(category);
let rows = std::fs::read_to_string(diagnostic_path.as_ref().unwrap()).unwrap();
assert_eq!(complete_original_occurrences(&rows, "InterfaceNotFound"), 2,
"both conflicting source originals must be complete before classification");
let public = request.dispatch_recorder().record_result::<()>(Err(category)).unwrap_err();
assert_eq!(public.code(), "saddle.diagnostic.source_unavailable");
drop(request);
return Err(public);
}
if matches!(mode, 19 | 20 | 21) {
let user_id = if mode == 19 { "" } else if mode == 20 {
"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
} else { "valid" };
let trace_id = "x".repeat(saddle_boundary::ingress::MAX_ID_BYTES + 8);
let source = crate::programming::ProfuseGwContext::from_framework_with_request(
user_id, if mode == 21 { &trace_id } else { "trace" },
"rpc", "zone", "idc", "env", Some(&request),
).err().expect("invalid context must fail at its creation point");
let rows = std::fs::read_to_string(diagnostic_path.as_ref().unwrap()).unwrap();
assert!(complete_original_with(&rows, &source.to_string()),
"context source must be written before leaving its constructor");
let public = request.dispatch_recorder().record_result::<()>(
Err(crate::programming::ProfuseGwDispatchError::ContextInvalid)).unwrap_err();
drop(request);
return Err(public);
}
if matches!(mode, 22 | 24) {
struct SaturatingRoute([u8; 32 * 1024]);
impl std::future::Future for SaturatingRoute {
type Output = ();
fn poll(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<()> {
std::hint::black_box(&self.0);
Poll::Pending
}
}
let memory = request.framework_future_memory();
let mut held = Vec::new();
let mut constructed = 0;
let error = loop {
assert!(held.len() < 1024, "finite framework reserve must reject a saturated route");
let before = constructed;
match saddle_admission::framework_route_future_at_source(
Some(&memory), || {
constructed += 1;
SaturatingRoute([0; 32 * 1024])
}, |error| request.record_framework_future_error_at_source(error)) {
Ok(future) => held.push(future),
Err(error) => {
assert_eq!(constructed, before, "failed reserve precedes route construction");
break error;
}
}
};
assert!(matches!(error, saddle_admission::AdmissionError::FrameworkReserveExceeded { .. }));
let rows = std::fs::read_to_string(diagnostic_path.as_ref().unwrap()).unwrap();
assert!(complete_original_with(&rows, "FrameworkReserveExceeded"),
"framework future reserve source survives early writer shutdown");
let public = request.dispatch_recorder().record_result::<()>(
Err(crate::programming::ProfuseGwDispatchError::FrameworkCapacityUnavailable)).unwrap_err();
assert_eq!(public.code(), "saddle.response.capacity_rejected");
drop(held);
drop(request);
return Err(public);
}
if mode == 23 {
let source = crate::programming::ProfuseGwDispatchError::InterfaceNotFound;
let source = request.record_dispatch_at_source(source);
let public = request.dispatch_recorder().record_result::<()>(Err(source)).unwrap_err();
assert_eq!(public.code(), "saddle.process.dispatch_failed");
drop(request);
return Err(public);
}
if matches!(mode, 25 | 26) {
request.record_dispatch_at_source(crate::programming::ProfuseGwDispatchError::InterfaceNotFound);
let recovered = request.dispatch_recorder().record_result(Ok(()));
drop(request);
recovered.unwrap();
return Ok(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}".to_vec());
}
if mode == 31 {
request.retain_test_parameter_failure("formal parameter failure without source output", false);
drop(request);
return Ok(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}".to_vec());
}
if mode == 33 {
request.retain_test_parameter_failure("formal confirmed slot failure recovered before business success", true);
drop(request);
return Ok(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}".to_vec());
}
if mode == 32 {
request.retain_test_recovered_scope_original(
"formal historical parameter failure recovered before business success"
).await;
drop(request);
return Ok(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}".to_vec());
}
if matches!(mode, 34 | 35) {
let failure = match request.parameter_json("{actual-invalid-json").await {
Ok(_) => panic!("invalid JSON must fail"),
Err(failure) => failure,
};
assert!(matches!(failure, saddle_db::internal::ParameterConstructionError::InvalidJson));
drop(request);
return Ok(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}".to_vec());
}
if mode == 40 {
let value = request.parameter_json("{\"a\":1}").await.unwrap();
let recorded = request.recorded_json(value, "response.payload").unwrap();
let encoding = request.response_encoding_context();
drop(request);
let response = crate::ingress::ProfuseGwResponse::<_, ResponseCapacityCode>::success((recorded,));
let wire = encoding.encode(move |memory|
crate::profusegw_http::encode_profusegw_http1_managed(&response, memory)).await?;
return Ok(wire.as_slice().to_vec());
}
if mode == 41 {
struct FailingWriter;
impl std::io::Write for FailingWriter {
fn write(&mut self, _: &[u8]) -> std::io::Result<usize> {
Err(std::io::ErrorKind::BrokenPipe.into())
}
fn flush(&mut self) -> std::io::Result<()> { Ok(()) }
}
let value = request.parameter_json("{\"a\":1}").await.unwrap();
let recorded = request.recorded_json(value, "response.payload").unwrap();
let encoding = request.response_encoding_context();
drop(request);
return encoding.encode(move |_memory| {
let error = serde_json::to_writer(FailingWriter, &recorded).unwrap_err();
source_barrier.store(1, Ordering::Release);
let until = std::time::Instant::now() + Duration::from_secs(3);
while source_barrier.load(Ordering::Acquire) != 2 {
assert!(std::time::Instant::now() < until, "independent source read before serializer returns");
std::hint::spin_loop();
}
Err(crate::profusegw_http::ManagedResponseEncodeError::Serialization(error))
}).await.map(|wire| wire.as_slice().to_vec());
}
drop(request);
if mode==8 { return Err(crate::process::process_response_capacity_error()); }
Ok(b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}".to_vec())
}
};
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let mut prepared = prepare::<(), (), _, _>("app", &lease, &dispatch).unwrap();
if mode == 14 {
fn fill<O: Send + 'static, T: Send + 'static>(
collection: &mut saddle_runtime::request_task::reserved_set::ReservedTaskCollection<O,T>,
remaining: usize,
) {
let stride = std::mem::size_of::<(
tokio::task::Id,
saddle_runtime::request_task::reserved::ReservedTaskTicket<O>,
)>();
collection.try_capacity((remaining - 64) / stride).unwrap();
}
let snapshot = lease.resource_snapshot().unwrap();
let remaining = snapshot.framework_capacity - snapshot.framework_charged;
fill(&mut prepared.recognizing, remaining);
assert!(lease.resource_snapshot().unwrap().framework_capacity
- lease.resource_snapshot().unwrap().framework_charged < 128);
}
let account_capacity=lease.resource_snapshot().unwrap().account_capacity;
let stop = Arc::new(AtomicBool::new(false));
let holder = if matches!(mode, 6 | 7 | 9 | 10) {Some(lease.try_ingress().unwrap())} else {None};
let total = lease.resource_snapshot().unwrap().managed_capacity;
let worker = tokio::spawn(run(
Some(plan_listener_input),
prepared,
listener,
saddle_boundary::ingress::ProfuseGwListenerAdapter::new("app").unwrap(),
(),
dispatch,
BusinessConfig::unit(),
None,
lease,
observer,
database_state,
factory_source,
stop.clone(),
));
if matches!(mode, 23 | 24 | 26 | 28 | 30) {
diagnostic.as_mut().unwrap().0.shutdown();
}
if mode == 11 {
tokio::time::sleep(Duration::from_millis(1300)).await;
let handshake = std::path::PathBuf::from(
std::env::var_os("SADDLE_I027_PRESSURE_HANDSHAKE").unwrap());
std::fs::write(handshake.join("ready"), b"ready").unwrap();
let until = std::time::Instant::now() + Duration::from_secs(7);
while !handshake.join("resume").exists() {
assert!(std::time::Instant::now() < until, "parent must resume paused listener");
tokio::time::sleep(Duration::from_millis(10)).await;
}
tokio::time::sleep(Duration::from_millis(300)).await;
}
let mut client = TcpStream::connect(address).await.unwrap();
let base=br#"{"requestData":{"label":"held","rows":[["copy"]]},"profuseGwContext":{"userInfo":{"userId":"u"},"traceInfo":{"rpcId":"0"},"ldcInfo":{"zone":"z","idc":"i","env":"test"}},"target":{"app":"app","interfaceId":"route"}}"#;
let body:std::borrow::Cow<'_,[u8]>=match mode {
36 | 38 | 39 => {
let value=match mode {36=>"\"invalid-integer\"",38=>"18446744073709551616",_=>"7"};
std::borrow::Cow::Owned(String::from_utf8(base.to_vec()).unwrap()
.replace("\"requestData\":{",&format!("\"requestData\":{{\"product_id\":{value},"))
.into_bytes())
}
37 => std::borrow::Cow::Owned(String::from_utf8(base.to_vec()).unwrap()
.replace("\"label\":\"held\"", "\"label\":]").into_bytes()),
_ => std::borrow::Cow::Borrowed(base),
};
let head = format!(
"POST /saddle/v1/ingress/profusegw/invoke HTTP/1.1\r\nContent-Type: application/json\r\nContent-Length: {}\r\nX-Request-Id: r\r\nX-Call-Id: c\r\n\r\n",
body.len()
);
let head = if mode == 1 {
let deadline = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis()
+ 300;
head.replace(
"\r\n\r\n",
&format!("\r\nX-Deadline-Unix-Ms: {deadline}\r\n\r\n"),
)
} else {
head
};
client.write_all(head.as_bytes()).await.unwrap();
client.write_all(&body).await.unwrap();
tokio::time::sleep(Duration::from_millis(150)).await;
if mode == 41 {
let until = std::time::Instant::now() + Duration::from_secs(3);
while source_barrier.load(Ordering::Acquire) != 1 {
assert!(std::time::Instant::now() < until, "serializer reached original-record barrier");
tokio::time::sleep(Duration::from_millis(5)).await;
}
let path = diagnostic.as_ref().unwrap().1.join("saddle.emergency.log");
let original = std::fs::read_to_string(path).unwrap();
let complete = complete_original_with(&original, "broken pipe");
let correct_field = original.contains("response.payload");
let correct_request = original.lines().filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "request_error_original" && row["channel"] == "context")
.filter_map(|row| serde_json::from_str::<serde_json::Value>(row["payload"].as_str()?).ok())
.any(|payload| payload["context"]["request"]["value"] == "r");
source_barrier.store(2, Ordering::Release);
assert!(complete && correct_field && correct_request,
"source before return: complete={complete}, field={correct_field}, request={correct_request}");
}
let mut queued=Vec::new();
if matches!(mode, 6 | 7 | 10) {
assert!(INPUT_PLANNED.load(Ordering::Acquire), "input must be planned before waiting");
let memory = holder.as_ref().unwrap().memory().unwrap();
let available = match memory.try_byte_buffer(total) {
Err(saddle_admission::AdmissionError::ProcessCapacityExceeded { available, .. }) => available,
_ => panic!("real remaining shared bytes must be below total after recognition"),
};
let held_bytes = if mode == 7 {
available.saturating_sub(crate::profusegw_http::RESPONSE_BUFFER_BYTES / 2)
} else { available };
let held = memory.try_byte_buffer(held_bytes).unwrap();
tokio::time::sleep(Duration::from_millis(1200)).await;
assert_eq!(entered.load(Ordering::SeqCst), if mode == 7 { 1 } else { 0 },
"only known input/execute storage gates initial dispatch");
drop(held);
}
let response_hold=if mode==9 {
tokio::time::timeout(Duration::from_secs(3),started.notified()).await.unwrap();
let memory=holder.as_ref().unwrap().memory().unwrap();
let available=match memory.try_byte_buffer(total+1) {
Err(saddle_admission::AdmissionError::BudgetExceeded { available, .. } |
saddle_admission::AdmissionError::ProcessCapacityExceeded { available, .. }) => available,
other => panic!("real response growth remaining bytes: {other:?}"),
};
let global=match memory.try_byte_buffer(available) {
Err(saddle_admission::AdmissionError::ProcessCapacityExceeded { available, .. }) => available,
Ok(probe) => {drop(probe);available},
other => panic!("real global remaining bytes: {other:?}"),
};
let held=memory.try_byte_buffer(available.min(global)-256).unwrap();
release.notify_one();
Some(held)
} else {None};
if mode==5 {
tokio::time::timeout(Duration::from_secs(3),started.notified()).await.unwrap();
for id in 1..=2 {
let mut peer=TcpStream::connect(address).await.unwrap();
let request=String::from_utf8(body.to_vec()).unwrap().replace("\"requestData\":{",&format!("\"requestData\":{{\"id\":{id},"));
let header=head.replace(&format!("Content-Length: {}",body.len()),&format!("Content-Length: {}",request.len()));
peer.write_all(header.as_bytes()).await.unwrap();
peer.write_all(request.as_bytes()).await.unwrap();
queued.push(peer);
tokio::time::sleep(Duration::from_millis(40)).await;
}
assert_eq!(entered.load(Ordering::SeqCst),3,"predependency business must exceed the one-connection DB pool");
release.notify_waiters();
}
if mode==12 {
let directory = &diagnostic.as_ref().unwrap().1;
let until = std::time::Instant::now() + Duration::from_secs(5);
loop {
let text = std::fs::read_to_string(directory.join("saddle.emergency.log")).unwrap_or_default();
if text.lines().filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.any(|row| row["event"] == "framework.execution_pressure"
&& row["state"] == "pressured" && row["valid"] == true) { break; }
assert!(std::time::Instant::now() < until, "three actual high windows must pressure the gate");
tokio::time::sleep(Duration::from_millis(20)).await;
}
let mut peer=TcpStream::connect(address).await.unwrap();
let request=String::from_utf8(body.to_vec()).unwrap()
.replace("\"requestData\":{","\"requestData\":{\"id\":1,");
let header=head.replace(&format!("Content-Length: {}",body.len()),&format!("Content-Length: {}",request.len()))
.replace("X-Request-Id: r\r\n","X-Request-Id: r2\r\n");
peer.write_all(header.as_bytes()).await.unwrap();
peer.write_all(request.as_bytes()).await.unwrap();
queued.push(peer);
tokio::time::sleep(Duration::from_millis(150)).await;
assert_eq!(entered.load(Ordering::SeqCst),1,"valid sustained CPU pressure must hold new business");
}
if mode==13 {
tokio::time::sleep(Duration::from_millis(3100)).await;
let mut peer=TcpStream::connect(address).await.unwrap();
let request=String::from_utf8(body.to_vec()).unwrap()
.replace("\"requestData\":{","\"requestData\":{\"id\":1,");
let header=head.replace(&format!("Content-Length: {}",body.len()),&format!("Content-Length: {}",request.len()))
.replace("X-Request-Id: r\r\n","X-Request-Id: r2\r\n");
peer.write_all(header.as_bytes()).await.unwrap();
peer.write_all(request.as_bytes()).await.unwrap();
queued.push(peer);
tokio::time::sleep(Duration::from_millis(150)).await;
assert_eq!(entered.load(Ordering::SeqCst),2,"I/O waiting alone must not close CPU gate");
release.notify_one();
}
if mode == 2 {
stop.store(true, Ordering::Release);
}
if mode == 3 {
drop(client);
tokio::time::sleep(Duration::from_millis(150)).await;
client = TcpStream::connect(address).await.unwrap();
client.write_all(head.as_bytes()).await.unwrap();
client.write_all(&body).await.unwrap();
}
if mode == 4 {
let mut held=Vec::new();
for _ in 1..account_capacity {
let socket=TcpStream::connect(address).await.unwrap();
held.push(socket);
}
tokio::time::sleep(Duration::from_millis(30)).await;
let mut overflow = TcpStream::connect(address).await.unwrap();
let mut rejected = Vec::new();
assert!(tokio::time::timeout(
Duration::from_millis(50), overflow.read_to_end(&mut rejected),
).await.is_err(), "initial account capacity must grow");
assert!(rejected.is_empty());
}
assert_eq!(
entered.load(Ordering::SeqCst),
if mode==5 {3} else if mode==13 {2} else if matches!(mode,7|9|12|41) {1} else {0},
"gate uses known initial demand rather than the maximum success wire"
);
let mut response = Vec::new();
tokio::time::timeout(Duration::from_secs(4), client.read_to_end(&mut response))
.await
.unwrap()
.unwrap();
if mode == 40 {
if let Some(path) = std::env::var_os("SADDLE_L_JSON_WIRE_EVIDENCE") {
std::fs::write(path, &response).unwrap();
}
}
match mode {
1 => assert!(
response.starts_with(b"HTTP/1.1 408"),
"{}",
String::from_utf8_lossy(&response)
),
2 => assert!(
response.is_empty(),
"shutdown must not invent a business response"
),
8 | 9 | 14 | 16 | 22 | 24 | 29 | 30 | 31 | 35 => assert!(response.starts_with(b"HTTP/1.1 503"),"{}",String::from_utf8_lossy(&response)),
36 | 37 | 38 => assert!(response.starts_with(b"HTTP/1.1 400"),"{}",String::from_utf8_lossy(&response)),
15 | 17 | 18 | 19 | 20 | 21 | 23 | 28 => assert!(response.starts_with(b"HTTP/1.1 500"),"{}",String::from_utf8_lossy(&response)),
41 => assert!(!response.starts_with(b"HTTP/1.1 200"), "faulting serializer cannot produce success"),
_ => assert!(
response.starts_with(b"HTTP/1.1 200"),
"{}",
String::from_utf8_lossy(&response)
),
}
if mode == 40 {
assert!(response.windows(b"{\"a\":1}".len()).any(|part| part == b"{\"a\":1}"),
"validated DB JSON must compose into the real response");
}
for mut peer in queued {
let mut response=Vec::new();
tokio::time::timeout(Duration::from_secs(if mode==12 {5} else {2}),peer.read_to_end(&mut response)).await.unwrap().unwrap();
assert!(response.starts_with(b"HTTP/1.1 200"),"{}",String::from_utf8_lossy(&response));
}
if mode==5 {assert_eq!(*order.lock().unwrap(),[0,1,2]);}
if mode==12 {assert_eq!(*order.lock().unwrap(),[0,1],"one effect per request after CPU recovery");}
if mode==13 {assert_eq!(*order.lock().unwrap(),[0,1],"I/O comparison keeps one effect per request");}
drop((response_hold,holder));
if mode==9 {
let mut peer=TcpStream::connect(address).await.unwrap();
let replay_body=String::from_utf8(body.to_vec()).unwrap().replace("\"requestData\":{","\"requestData\":{\"id\":1,");
let replay_head=head.replace(&format!("Content-Length: {}",body.len()),&format!("Content-Length: {}",replay_body.len()))
.replace("X-Request-Id: r\r\n","X-Request-Id: r2\r\n");
peer.write_all(replay_head.as_bytes()).await.unwrap();
peer.write_all(replay_body.as_bytes()).await.unwrap();
release.notify_one();
let mut recovered=Vec::new();
tokio::time::timeout(Duration::from_secs(4),peer.read_to_end(&mut recovered)).await.unwrap().unwrap();
assert!(recovered.starts_with(b"HTTP/1.1 200"),"{}",String::from_utf8_lossy(&recovered));
assert_eq!(*order.lock().unwrap(),[0,1],"failed encoding keeps only the original effect");
}
assert_eq!(
entered.load(Ordering::SeqCst),
if matches!(mode, 1 | 2 | 14 | 36 | 37 | 38) { 0 } else if mode==5 {3} else if matches!(mode,9|12|13) {2} else { 1 }
);
stop.store(true, Ordering::Release);
worker.await.unwrap();
drop(diagnostic_handle);
if let Some((mut writer, directory)) = diagnostic.take() {
let until = std::time::Instant::now() + Duration::from_secs(3);
loop {
if writer.shutdown() == saddle_observability::DiagnosticShutdown::Finished { break; }
assert!(std::time::Instant::now() < until, "emergency writer drains original FIFO facts");
tokio::time::sleep(Duration::from_millis(10)).await;
}
let text = std::fs::read_to_string(directory.join("saddle.emergency.log")).unwrap();
if matches!(mode,36|37|38) {
let marker=match mode {36=>"invalid-integer",37=>"expected value",_=>"18446744073709551616"};
assert!(complete_original_with(&text, marker), "actual input failure original must be complete");
assert!(text.contains("request r"), "validated request identity");
if mode != 37 {
assert!(text.contains("requestData.product_id"));
assert!(text.contains("u64"));
}
assert!(!text.contains("outbound-only"), "credentials never enter originals");
} else if mode == 14 {
let originals: Vec<serde_json::Value> = text.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "request_error_original")
.collect();
assert!(originals.iter().any(|row| {
if row["channel"] != "debug" { return false; }
let Some(payload) = row["payload"].as_str() else { return false; };
let Some((_, rest)) = payload.split_once("Storage(FrameworkReserveExceeded { requested: ") else { return false; };
let Some((requested, rest)) = rest.split_once(", available: ") else { return false; };
let Some((available, _)) = rest.split_once('}') else { return false; };
requested.trim().parse::<usize>().ok().zip(available.trim().parse::<usize>().ok())
.is_some_and(|(requested, available)| requested > available)
}), "original same-lock rejection shortage is recorded before fixed fallback");
let resource: Vec<serde_json::Value> = text.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "framework.ingress_resource_decision")
.collect();
assert!(resource.iter().any(|row|
row["decision"] == "capacity_rejected"
&& row["resource_source"] == "admission_ledger_same_lock"
&& row["resource_valid"] == true
&& row["decision_monotonic_ms"].is_null()
&& row["decision_time_valid"] == false
&& row["decision_time_source"] == "not_captured_by_ledger_transaction"
&& row["timestamp_role"] == "output_enqueued"
&& row["resource"]["domain"] == "framework"
&& row["resource"]["requested_bytes"].as_u64().unwrap()
> row["resource"]["available_bytes"].as_u64().unwrap()
&& row["resource"]["available_bytes"].as_u64().unwrap() < 128
&& row["context"]["request"]["state"] == "not_established"),
"terminal pre-root 503 must expose original same-lock byte decision");
println!("I027_FORMAL_FRAMEWORK_REJECTION PASS {}", directory.display());
} else if mode == 13 {
let rows: Vec<serde_json::Value> = text.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "framework.execution_pressure")
.collect();
assert!(rows.iter().any(|row|
row["state"] == "open" && row["valid"] == true
&& row["process_cpu_ratio"].as_f64().unwrap() <= 0.65));
assert!(!rows.iter().any(|row| row["state"] == "pressured"),
"elapsed I/O wait is not sustained CPU pressure");
println!("I027_FORMAL_IO_CONTROL PASS {}", directory.display());
} else if mode == 12 {
let rows: Vec<serde_json::Value> = text.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "framework.execution_pressure")
.collect();
let pressured = rows.iter().position(|row| row["state"] == "pressured"
&& row["valid"] == true
&& row["process_cpu_ratio"].as_f64().unwrap() >= 0.85)
.expect("real sustained process CPU pressure");
assert!(rows.iter().skip(pressured + 1).any(|row|
row["state"] == "open" && row["valid"] == true),
"three real low windows recover the original gate");
println!("I027_FORMAL_CPU_PRESSURE PASS {}", directory.display());
} else if mode == 11 {
let rows: Vec<serde_json::Value> = text.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "framework.execution_pressure")
.collect();
let valid_at = rows.iter().find(|row|
row["window_valid"] == true && row["valid"] == true)
.and_then(|row| row["observed_monotonic_ms"].as_u64())
.expect("real valid window before process suspension");
assert!(rows.iter().any(|row|
row["source"] == "process_cpu_clock+tokio_workers+sampler_delay"
&& row["original_error"] == "pressure.missing_window"
&& row["window_valid"] == false
&& row["valid"] == false
&& row["state"] == "invalid"
&& row["observed_monotonic_ms"].as_u64().unwrap() > valid_at
&& row["scheduling_delay_ms"].as_u64().unwrap() >= 2000
&& row["decision_monotonic_ms"].as_u64().unwrap()
>= row["observed_monotonic_ms"].as_u64().unwrap()));
println!("I027_FORMAL_PRESSURE_INVALID PASS {}", directory.display());
} else if mode == 16 {
assert!(!complete_original_with(&text, "InterfaceNotFound"),
"missing receipt cannot issue a written-source claim");
assert!(!complete_original_with(&text, "FrameworkReserveExceeded"));
} else if matches!(mode, 31 | 35) {
assert!(!complete_original_with(&text, if mode == 35 {
"actual-invalid-json"
} else { "formal parameter failure without source output" }));
assert!(text.lines().filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.any(|row| row["event"] == "request_failure_boundary"
&& row["original_capture"] != "complete_written"
&& (mode != 35 ||
(row["source_submission"] == "output_unavailable"
&& row["occurrence"]["diagnostic_id"].as_u64().is_some()
&& row["axes"]["axes"]["delivery"] == "local_write_complete"))),
"unconfirmed source keeps its negative terminal fact");
} else if matches!(mode, 15 | 17 | 18 | 19 | 20 | 21 | 22 | 23 | 24 | 25 | 26 | 28 | 29 | 30 | 32 | 33 | 34) {
let cause = match mode {
34 => "actual-invalid-json",
32 => "formal historical parameter failure recovered before business success",
33 => "formal confirmed slot failure recovered before business success",
17 => "external decode cause",
18 => "IdentityMismatch",
19 => "MissingUserId",
20 => "UserIdTooLong",
21 => "ContextTooLong",
22 | 24 => "FrameworkReserveExceeded",
28 => "closed writer long external cause",
_ => "InterfaceNotFound",
};
assert!(complete_original_with(&text, cause));
if matches!(mode, 32 | 33 | 34) {
let records: Vec<serde_json::Value> = text.lines()
.filter_map(|line| serde_json::from_str(line).ok()).collect();
if mode == 34 {
assert!(records.iter().any(|row|
row["event"] == "request_error_original"
&& row["payload"].as_str().is_some_and(|value|
value.contains("InvalidJson"))),
"the actual admission error remains in the exposed chain");
}
let originals: Vec<_> = records.iter()
.filter(|row| row["event"] == "request_error_original"
&& row["channel"] == "description"
&& row["payload"].as_str().is_some_and(|value| value.contains(cause)))
.collect();
assert_eq!(originals.len(), 1, "one retained source original");
assert!(records.iter().any(|row|
row["event"] == "request_failure_boundary"
&& row["occurrence"] == originals[0]["occurrence"]
&& row["source_submission"] == "written"
&& row["original_capture"] == "complete_written"
&& row["axes"]["axes"]["delivery"] == "local_write_complete"),
"the same source reaches terminal after actual delivery");
}
if matches!(mode, 29 | 30) {
assert_eq!(complete_original_occurrences(&text, cause), 2);
}
} else if mode == 40 {
assert!(!complete_original_with(&text, "db.json.serialize"),
"valid composed JSON has no serializer source error");
} else if mode == 41 {
assert!(complete_original_with(&text, "broken pipe"));
} else {
let queued: Vec<serde_json::Value> = text.lines()
.filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
.filter(|row| row["event"] == "framework.capacity"
&& row["construction"]["status"] == "queued_for_retry")
.collect();
assert!(queued.len() >= 4 && queued.len() % 4 == 0,
"every actual FIFO retry emits its original four dimensions");
for decision in queued.chunks_exact(4) {
let resource = &decision[0]["resource"];
assert!(resource["managed_demand_bytes"].as_u64().unwrap()
> resource["managed_available_bytes"].as_u64().unwrap());
assert_eq!(decision.iter().map(|row|row["capacity_dimension"].as_str().unwrap()).collect::<Vec<_>>(),
["cpu", "memory", "database", "profuse_contract"]);
assert!(decision.iter().all(|row|
&row["resource"] == resource
&& row["resource_source"] == "admission_ledger_same_lock"
&& row["context"]["request"]["state"] == "not_established"
&& row["context"]["route"]["state"] == "not_established"
&& row["recognized_ingress"]["request"] == "r"
&& row["recognized_ingress"]["route"] == "route"
&& row["context"]["scope"]["state"] == "not_established"));
}
println!("I027_WAIT_FACTS PASS {}", directory.display());
}
}
if let Some(entry) = database_entry {
saddle_runtime::ComponentLifecycle::shutdown(&entry).await.unwrap();
}
assert!(
std::process::Command::new("kill")
.args(["-TERM", &std::process::id().to_string()])
.status()
.unwrap()
.success()
);
Ok(saddle_runtime::Application::new())
},
))
}));
let result = coordinator.as_mut().poll(&mut Context::from_waker(Waker::noop()));
match result {
Poll::Ready(Ok(())) => {},
Poll::Ready(Err(saddle_runtime::profusegw::ProfuseGwCoordinatorFailure::Lifecycle(error))) =>
panic!("mode {mode}: lifecycle {}", error.code()),
Poll::Ready(Err(saddle_runtime::profusegw::ProfuseGwCoordinatorFailure::Finalization { error, .. })) =>
panic!("mode {mode}: finalization {error:?}"),
Poll::Ready(Err(saddle_runtime::profusegw::ProfuseGwCoordinatorFailure::Startup(_))) =>
panic!("mode {mode}: startup failure"),
Poll::Ready(Err(saddle_runtime::profusegw::ProfuseGwCoordinatorFailure::LifecycleWithCleanup { primary, .. })) =>
panic!("mode {mode}: lifecycle with cleanup {}", primary.code()),
Poll::Ready(Err(saddle_runtime::profusegw::ProfuseGwCoordinatorFailure::OwnerMissing { error, .. })) =>
panic!("mode {mode}: owner missing {}", error.code()),
Poll::Ready(Err(_)) => panic!("mode {mode}: other coordinator failure"),
Poll::Pending => panic!("mode {mode}: coordinator pending"),
}
}
#[test]
fn formal_listener_recognizes_before_cold_gate_and_releases_original_account() {
formal_listener_case(0);
}
#[test]
fn formal_listener_input_decode_originals_and_normal_product_id() {
const CHILD: &str = "SADDLE_RD_INPUT_MODE";
if let Ok(mode) = std::env::var(CHILD) {
formal_listener_case(mode.parse().unwrap());
return;
}
for mode in [36, 37, 38, 39] {
assert!(std::process::Command::new(std::env::current_exe().unwrap())
.args(["--exact", "process::reserved_entry::lifecycle_tests::formal_listener_input_decode_originals_and_normal_product_id"])
.env(CHILD, mode.to_string()).status().unwrap().success(), "mode {mode}");
}
}
#[test]
fn formal_listener_waiting_keeps_upstream_deadline() {
formal_listener_case(1);
}
#[test]
fn formal_listener_shutdown_drains_waiting_without_business() {
formal_listener_case(2);
}
#[test]
fn formal_listener_disconnect_releases_waiting_slot() {
formal_listener_case(3);
}
#[test]
fn formal_listener_crosses_initial_account_backing_without_rejection() {
formal_listener_case(4);
}
#[test]
fn formal_listener_framework_exhaustion_still_delivers_fixed_503() {
formal_listener_case(14);
}
#[test]
fn formal_listener_predependency_business_exceeds_database_pool() { formal_listener_case(5); }
#[test]
fn formal_listener_input_peak_waits_before_dispatch_and_recovers_on_release() { formal_listener_case(6); }
#[test]
fn formal_listener_waiting_decision_reads_original_diagnostic() {
let status = std::process::Command::new(std::env::current_exe().unwrap())
.arg("formal_listener_waiting_decision_child")
.arg("--ignored")
.env("SADDLE_I027_WAIT_FACTS_CHILD", "1")
.status().unwrap();
assert!(status.success(), "isolated formal listener diagnostic readback");
}
#[test]
#[ignore = "run by the isolated parent test"]
fn formal_listener_waiting_decision_child() {
if std::env::var_os("SADDLE_I027_WAIT_FACTS_CHILD").is_some() {
formal_listener_case(10);
}
}
#[test]
fn formal_listener_real_missing_pressure_window_reads_original_diagnostic() {
let handshake = std::env::temp_dir().join(format!("saddle-i027-pressure-handshake-{}", std::process::id()));
std::fs::create_dir_all(&handshake).unwrap();
let mut child = std::process::Command::new(std::env::current_exe().unwrap())
.arg("formal_listener_pressure_invalid_child")
.arg("--ignored")
.env("SADDLE_I027_PRESSURE_CHILD", "1")
.env("SADDLE_I027_PRESSURE_HANDSHAKE", &handshake)
.spawn().unwrap();
let pid = child.id().to_string();
let mut stopped = false;
let run = (|| -> std::result::Result<(), &'static str> {
let until = std::time::Instant::now() + Duration::from_secs(6);
while !handshake.join("ready").exists() {
if child.try_wait().map_err(|_| "child wait failed")?.is_some() { return Err("child exited before valid window"); }
if std::time::Instant::now() >= until { return Err("listener did not establish valid window"); }
std::thread::sleep(Duration::from_millis(10));
}
if !std::process::Command::new("kill").args(["-STOP", &pid]).status()
.map_err(|_| "SIGSTOP unavailable")?.success() { return Err("SIGSTOP failed"); }
stopped = true;
std::thread::sleep(Duration::from_millis(2200));
if !std::process::Command::new("kill").args(["-CONT", &pid]).status()
.map_err(|_| "SIGCONT unavailable")?.success() { return Err("SIGCONT failed"); }
stopped = false;
std::fs::write(handshake.join("resume"), b"resume").map_err(|_| "resume marker failed")?;
Ok(())
})();
if stopped { let _ = std::process::Command::new("kill").args(["-CONT", &pid]).status(); }
if run.is_err() { let _ = child.kill(); }
let until = std::time::Instant::now() + Duration::from_secs(10);
let status = loop {
if let Some(status) = child.try_wait().unwrap() { break status; }
if std::time::Instant::now() >= until {
let _ = child.kill();
break child.wait().unwrap();
}
std::thread::sleep(Duration::from_millis(10));
};
assert!(run.is_ok() && status.success(), "formal pressure diagnostic: {run:?}, {status:?}");
}
#[test]
#[ignore = "run by the isolated parent test"]
fn formal_listener_pressure_invalid_child() {
if std::env::var_os("SADDLE_I027_PRESSURE_CHILD").is_some() {
formal_listener_case(11);
}
}
#[test]
fn formal_listener_real_cpu_pressure_holds_and_recovers_original_fifo() {
let status = std::process::Command::new(std::env::current_exe().unwrap())
.arg("formal_listener_cpu_pressure_child")
.arg("--ignored")
.env("SADDLE_I027_CPU_CHILD", "1")
.status().unwrap();
assert!(status.success(), "isolated finite CPU pressure journey");
}
#[test]
#[ignore = "run by the isolated parent test"]
fn formal_listener_cpu_pressure_child() {
if std::env::var_os("SADDLE_I027_CPU_CHILD").is_some() {
formal_listener_case(12);
}
}
#[test]
fn formal_listener_io_wait_keeps_cpu_gate_open_for_same_dispatch() {
let status = std::process::Command::new(std::env::current_exe().unwrap())
.arg("formal_listener_io_control_child")
.arg("--ignored")
.env("SADDLE_I027_IO_CHILD", "1")
.status().unwrap();
assert!(status.success(), "isolated I/O wait control journey");
}
#[test]
#[ignore = "run by the isolated parent test"]
fn formal_listener_io_control_child() {
if std::env::var_os("SADDLE_I027_IO_CHILD").is_some() {
formal_listener_case(13);
}
}
#[test]
fn formal_listener_does_not_precommit_maximum_success_wire() { formal_listener_case(7); }
#[test]
fn formal_listener_response_capacity_failure_keeps_original_root() { formal_listener_case(8); }
#[test]
fn formal_listener_retained_parameter_failure_blocks_business_success() { formal_listener_case(31); }
#[test]
fn formal_listener_written_historical_failure_allows_business_recovery() { formal_listener_case(32); }
#[test]
fn formal_listener_confirmed_slot_history_allows_business_recovery() { formal_listener_case(33); }
#[test]
#[ignore = "requires a private MariaDB fixture and SADDLE_I_DB_CONNECTION_FILE"]
fn formal_listener_real_managed_parameter_failure_allows_recovery() {
if std::env::var_os("SADDLE_I_DB_CONNECTION_FILE").is_none() { return; }
for mode in [34, 35] {
let status = std::process::Command::new(std::env::current_exe().unwrap())
.arg("formal_listener_real_managed_parameter_child")
.arg("--ignored")
.env("SADDLE_I_REAL_PARAMETER_MODE", mode.to_string())
.status().unwrap();
assert!(status.success(), "isolated managed DB listener parameter failure mode {mode}");
}
}
#[test]
#[ignore = "requires a private MariaDB fixture and SADDLE_I_DB_CONNECTION_FILE"]
fn formal_listener_real_managed_json_response_composition() {
if std::env::var_os("SADDLE_I_DB_CONNECTION_FILE").is_none() { return; }
let status = std::process::Command::new(std::env::current_exe().unwrap())
.arg("formal_listener_real_managed_parameter_child")
.arg("--ignored")
.env("SADDLE_I_REAL_PARAMETER_MODE", "40")
.status().unwrap();
assert!(status.success(), "isolated managed JSON response composition");
}
#[test]
#[ignore = "requires a private MariaDB fixture and SADDLE_I_DB_CONNECTION_FILE"]
fn formal_listener_real_managed_json_serializer_source_before_return() {
if std::env::var_os("SADDLE_I_DB_CONNECTION_FILE").is_none() { return; }
let status = std::process::Command::new(std::env::current_exe().unwrap())
.arg("formal_listener_real_managed_parameter_child")
.arg("--ignored")
.env("SADDLE_I_REAL_PARAMETER_MODE", "41")
.status().unwrap();
assert!(status.success(), "isolated managed JSON serializer source before return: {status:?}");
}
#[test]
#[ignore = "run by the isolated parent test with a private MariaDB fixture"]
fn formal_listener_real_managed_parameter_child() {
if std::env::var_os("SADDLE_I_DB_CONNECTION_FILE").is_some() {
formal_listener_case(std::env::var("SADDLE_I_REAL_PARAMETER_MODE").unwrap().parse().unwrap());
}
}
#[test]
#[ignore = "requires a private MariaDB fixture and SADDLE_I_DB_CONNECTION_FILE"]
fn application_process_entry_real_parameter_failure() {
const CHILD: &str = "SADDLE_I_APPLICATION_PARAMETER_CHILD";
if std::env::var_os("SADDLE_I_DB_CONNECTION_FILE").is_none() { return; }
if let Some(root) = std::env::var_os(CHILD) {
use saddle_runtime::startup_assembly::StartupDbPoolFactory as _;
let root = std::path::PathBuf::from(root);
let output = saddle_observability::EmergencyDiagnostics::start(
&saddle_observability::FileLoggingConfig::new(&root,
saddle_observability::Rotation::Daily)).unwrap();
let handle = output.handle();
let connection_file = std::path::PathBuf::from(
std::env::var("SADDLE_I_DB_CONNECTION_FILE").unwrap());
let connection: serde_json::Value = serde_json::from_slice(
&std::fs::read(&connection_file).unwrap()).unwrap();
let mappings = connection_file.parent().unwrap().join("empty-mappings");
std::fs::create_dir_all(&mappings).unwrap();
let config = saddle_db::DatabaseConfig::new(connection["url"].as_str().unwrap())
.max_connections(1).name_mapping_directory(mappings);
let no_db_output = std::env::var_os("SADDLE_I_APPLICATION_NO_DB_OUTPUT").is_some();
let capacity_fault = std::env::var_os("SADDLE_I_APPLICATION_CAPACITY_FAULT").is_some();
let response = Arc::new(std::sync::Mutex::new(None));
let completed = response.clone();
let budget = budget_with_timeout(if capacity_fault { 30_000 } else { 5000 });
let result = poll_ready(coordinate_profusegw_app_run(budget, move |process| {
let (terminal, exit) = run_profusegw_owned_application_with_diagnostics(
process, output, move |lease| async move {
let startup = lease.take_database_startup_half().unwrap();
let observer = saddle_observability::Observer::with_writer(
Default::default(), std::io::sink()).unwrap();
let factory = saddle_db::internal::StartupManagedDatabaseFactory::new(
Some(config), observer.clone());
let factory = if no_db_output { factory } else {
factory.with_diagnostic_output(handle.clone())
};
let owner = factory.construct(saddle_admission::DbCreditProfile {
connections: 1, operations: 1,
}).await.unwrap();
let (db_entry, db_state) = ManagedDatabaseEntry::new(owner, startup,
"app".to_owned(), Some(handle.clone()));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let dispatch = move |_: saddle_boundary::ingress::AcceptedIngress, _: (),
_: BusinessConfig<()>, mut request: crate::database_capability::DatabaseRequest| async move {
let error = if capacity_fault {
let input = format!("capacity-source-start:{}", "X".repeat(34 * 1024 * 1024));
match request.parameter_text(&input).await {
Ok(_) => panic!("actual capacity construction must fail"),
Err(error) => error,
}
} else {
match request.parameter_json("{application-invalid-json").await {
Ok(_) => panic!("actual JSON construction must fail"),
Err(error) => error,
}
};
if capacity_fault {
match error {
saddle_db::internal::ParameterConstructionError::Memory(
saddle_admission::AdmissionError::BudgetExceeded { requested, available }
| saddle_admission::AdmissionError::ProcessCapacityExceeded { requested, available }
) => assert!(requested > available, "actual same-ledger byte shortage"),
_ => panic!("expected actual parameter memory shortage"),
}
} else {
assert!(matches!(error,
saddle_db::internal::ParameterConstructionError::InvalidJson));
}
Ok::<_, saddle_core::SaddleError>(
b"HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\n{}".to_vec())
};
let entry = ProcessEntry::new(Some(plan_listener_input), listener,
saddle_boundary::ingress::ProfuseGwListenerAdapter::new("app").unwrap(),
(), dispatch, BusinessConfig::unit(), None, lease, observer.clone(),
Some(db_state), Some(handle.clone()));
let mut app = saddle_runtime::Application::new();
app.install_lifecycle_observer_with_output(observer, "app", Some(handle));
app.register(db_entry)?;
app.register(entry)?;
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(100)).await;
let mut client = TcpStream::connect(address).await.unwrap();
let body = br#"{"requestData":{"label":"held","rows":[["copy"]]},"profuseGwContext":{"userInfo":{"userId":"u"},"traceInfo":{"rpcId":"0"},"ldcInfo":{"zone":"z","idc":"i","env":"test"}},"target":{"app":"app","interfaceId":"route"}}"#;
let head = format!("POST /saddle/v1/ingress/profusegw/invoke HTTP/1.1\r\nContent-Type: application/json\r\nContent-Length: {}\r\nX-Request-Id: r\r\nX-Call-Id: c\r\n\r\n", body.len());
client.write_all(head.as_bytes()).await.unwrap();
client.write_all(body).await.unwrap();
let mut received = Vec::new();
tokio::time::timeout(Duration::from_secs(if capacity_fault { 20 } else { 4 }),
client.read_to_end(&mut received))
.await.unwrap().unwrap();
*completed.lock().unwrap() = Some(received);
assert!(std::process::Command::new("kill")
.args(["-TERM", &std::process::id().to_string()])
.status().unwrap().success());
});
Ok(app)
});
assert_eq!(exit.shutdown, saddle_observability::DiagnosticShutdown::Finished);
std::future::ready(terminal)
}));
assert!(result.is_ok(), "real Application finalization");
let response = response.lock().unwrap();
let response = response.as_ref().expect("real HTTP response before shutdown");
assert!(response.starts_with(if no_db_output { b"HTTP/1.1 503" } else { b"HTTP/1.1 200" }),
"{}", String::from_utf8_lossy(response));
return;
}
for (no_db_output, capacity_fault) in [(false, false), (true, false), (false, true)] {
let root = std::env::temp_dir().join(format!("saddle-i-app-parameter-{}-{no_db_output}-{capacity_fault}",
std::process::id()));
std::fs::create_dir(&root).unwrap();
let mut command = std::process::Command::new(std::env::current_exe().unwrap());
command.arg("--exact")
.arg("process::reserved_entry::lifecycle_tests::application_process_entry_real_parameter_failure")
.arg("--ignored")
.env(CHILD, &root);
if no_db_output { command.env("SADDLE_I_APPLICATION_NO_DB_OUTPUT", "1"); }
if capacity_fault { command.env("SADDLE_I_APPLICATION_CAPACITY_FAULT", "1"); }
assert!(command.status().unwrap().success());
let rows = std::fs::read_to_string(root.join("saddle.emergency.log")).unwrap();
let records: Vec<serde_json::Value> = rows.lines()
.map(|line| serde_json::from_str(line).unwrap()).collect();
let boundary = records.iter().find(|row|
row["event"] == "request_failure_boundary")
.expect("actual Application request terminal");
let id = &boundary["occurrence"]["diagnostic_id"];
assert!(id.as_u64().is_some());
assert_eq!(boundary["axes"]["axes"]["delivery"], "local_write_complete");
assert_eq!(boundary["context"]["application"]["value"], "app");
assert_eq!(boundary["classification"], "db.parameters.construct");
assert_eq!(boundary["source_context"]["db_operation"]["value"],
if capacity_fault { "db.parameters.text" } else { "db.parameters.json" });
if no_db_output {
assert!(!complete_original_with(&rows, "application-invalid-json"));
assert_eq!(boundary["source_submission"], "output_unavailable");
assert_ne!(boundary["original_capture"], "complete_written");
} else {
let marker = if capacity_fault { "capacity-source-start" }
else { "application-invalid-json" };
assert!(complete_original_with(&rows, marker));
assert_eq!(complete_original_occurrences(&rows, marker), 1);
assert_eq!(boundary["source_submission"], "written");
assert_eq!(boundary["original_capture"], "complete_written");
let contexts: Vec<_> = records.iter().filter(|row|
row["event"] == "request_error_original"
&& row["channel"] == "context"
&& &row["occurrence"]["diagnostic_id"] == id).collect();
assert_eq!(contexts.len(), 1, "one original reaches Application finalization");
let context: serde_json::Value = serde_json::from_str(
contexts[0]["payload"].as_str().unwrap()).unwrap();
assert_eq!(context["context"]["application"]["value"], "app");
if capacity_fault {
let description = records.iter().filter(|row|
row["event"] == "request_error_original"
&& row["channel"] == "description" && row["cause_depth"] == 0
&& &row["occurrence"]["diagnostic_id"] == id)
.map(|row| row["payload"].as_str().unwrap()).collect::<String>();
assert!(description.contains("capacity-source-start"));
assert!(description.contains("BudgetExceeded { requested:")
|| description.contains("ProcessCapacityExceeded { requested:"));
assert!(description.len() >= 34 * 1024 * 1024,
"every byte of the failed parameter input survives source recording");
}
}
std::fs::remove_file(root.join("saddle.emergency.log")).unwrap();
std::fs::remove_dir(root).unwrap();
}
}
#[test]
fn formal_listener_actual_encoder_growth_failure_returns_503() { formal_listener_case(9); }
#[test]
fn formal_dispatch_source_is_written_before_failure_and_cleanup() {
for mode in [15, 16, 17, 18, 19, 20, 21, 22, 23, 24, 25, 26, 28, 29, 30] {
let status = std::process::Command::new(std::env::current_exe().unwrap())
.arg("formal_dispatch_written_child")
.arg("--ignored")
.env("SADDLE_D_WRITTEN_CHILD", mode.to_string())
.status().unwrap();
assert!(status.success(), "isolated formal dispatch written source mode {mode}");
}
}
#[test]
#[ignore = "run by the isolated parent test"]
fn formal_dispatch_written_child() {
if let Ok(mode) = std::env::var("SADDLE_D_WRITTEN_CHILD") {
formal_listener_case(mode.parse().unwrap());
}
}
#[test]
fn formal_response_wire_grows_with_result_and_refunds_after_failure() {
use crate::profusegw_http::{encode_profusegw_http1_managed, ManagedResponseEncodeError, RESPONSE_BUFFER_BYTES};
use crate::ingress::{ProfuseGwCode, ProfuseGwResponse};
enum Code { Unused }
impl ProfuseGwCode for Code {
const REGISTERED_CODES: &'static [&'static str] = &["UNUSED"];
fn stable_code(&self) -> &'static str { "UNUSED" }
}
let process = saddle_admission::prepare_profusegw_lightweight_profile(budget()).unwrap_or_else(|_|panic!("valid profile"));
let ingress = process.verified_profile().try_ingress().unwrap();
let memory = ingress.memory().unwrap();
let baseline = process.resource_snapshot().charged;
let small = ProfuseGwResponse::<_, Code>::success(String::from("small"));
let small_wire = memory.with_response_encoding(|| encode_profusegw_http1_managed(&small, &memory)).unwrap();
let small_charge = process.resource_snapshot().charged - baseline;
assert!(small_charge < RESPONSE_BUFFER_BYTES / 100, "small result cannot reserve maximum wire");
assert!(small_wire.as_slice().ends_with(b"\"small\"}"));
drop(small_wire);
assert_eq!(process.resource_snapshot().charged, baseline);
let large = ProfuseGwResponse::<_, Code>::success("a".repeat(16_384));
let large_wire = memory.with_response_encoding(|| encode_profusegw_http1_managed(&large, &memory)).unwrap();
let large_charge = process.resource_snapshot().charged - baseline;
assert!(large_charge > small_charge);
drop(large_wire);
assert_eq!(process.resource_snapshot().charged, baseline);
let total = process.resource_snapshot().managed_capacity;
let available = match memory.try_byte_buffer(total + 1) {
Err(saddle_admission::AdmissionError::ProcessCapacityExceeded { available, .. }) => available,
Err(saddle_admission::AdmissionError::BudgetExceeded { available, .. }) => available,
other => panic!("expected true remaining byte fact: {other:?}"),
};
let held = memory.try_byte_buffer(available - 256).unwrap();
let before = process.resource_snapshot().charged;
let rejected = memory.with_response_encoding(|| encode_profusegw_http1_managed(&large, &memory));
assert!(matches!(rejected, Err(ManagedResponseEncodeError::Admission(_))));
assert_eq!(process.resource_snapshot().charged, before, "failed growth refunds its old and new blocks");
drop(held);
let recovered = memory.with_response_encoding(|| encode_profusegw_http1_managed(&large, &memory)).unwrap();
drop(recovered);
drop((small,large,memory,ingress));
assert_eq!(process.resource_snapshot().charged, 0);
process.finish().unwrap();
}