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 {
let pending = saddle_admission::freeze_deployment_resource_budget(
1, 32, 5000, 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(5000),
))
.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>>>;
#[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,
diagnostic_handle: None,
}),
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>, label:Option<String>, rows:Option<Vec<Vec<String>>> }
impl saddle_admission::DecodeInput for ListenerInput {
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<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")?,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];
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) -> std::result::Result<usize,saddle_admission::AdmissionError> {
let peak = crate::programming::plan_accepted_profusegw::<ListenerInput>(input)?;
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) {
let directory = std::env::temp_dir().join(format!(
"saddle-i027-entry-facts-{}", std::process::id()));
std::fs::create_dir_all(&directory).unwrap();
let writer = saddle_observability::EmergencyDiagnostics::start(
&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 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 dispatch = move |input:saddle_boundary::ingress::AcceptedIngress, _:(), _:BusinessConfig<()>, request:crate::database_capability::DatabaseRequest| {
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 {
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());
}
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>,
)>();
// StorageDemand also bills the physical permit header.
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);
}
// Initial backing is not an ingress-count limit.
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,
None,
diagnostic_handle.clone(),
stop.clone(),
));
if mode == 11 {
// The parent suspends this whole process after the listener
// has produced a real valid window, then resumes it.
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 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()
);
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;
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();
// Beyond the normal cold window, input still requires its
// known peak; success wire is no longer prepaid at max.
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) {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();
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 => assert!(response.starts_with(b"HTTP/1.1 503"),"{}",String::from_utf8_lossy(&response)),
_ => assert!(
response.starts_with(b"HTTP/1.1 200"),
"{}",
String::from_utf8_lossy(&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) { 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 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 {
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());
}
}
drop(startup);
assert!(
std::process::Command::new("kill")
.args(["-TERM", &std::process::id().to_string()])
.status()
.unwrap()
.success()
);
Ok(saddle_runtime::Application::new())
},
))
}));
assert!(matches!(
coordinator
.as_mut()
.poll(&mut Context::from_waker(Waker::noop())),
Poll::Ready(Ok(()))
));
}
#[test]
fn formal_listener_recognizes_before_cold_gate_and_releases_original_account() {
formal_listener_case(0);
}
#[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_actual_encoder_growth_failure_returns_503() { formal_listener_case(9); }
#[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();
}