use super::*;
use crate::request_task::reserved::{
ReservedBorrowedFuture, ReservedRequestFailure, ReservedTaskContext,
};
use saddle_core::request_context::ContextLabel;
use std::{
alloc::Layout,
sync::atomic::{AtomicBool, Ordering},
};
#[derive(Clone, Copy, Debug)]
enum Case {
Ready,
PrePoll,
Pending,
Withdraw,
}
struct Input {
case: Case,
entered: Arc<AtomicBool>,
observer: Observer,
}
type Owner = ReservedDispatchOwner<Input>;
fn body(
owner: &mut Owner,
mut context: ReservedTaskContext,
) -> impl Future<Output = Result<u32, ReservedRequestFailure>> + Send + '_ {
async move {
owner.1.entered.store(true, Ordering::SeqCst);
assert_eq!(
owner.2.as_ref().unwrap().receipt.decision(),
ProfuseGwCapacityDecision::Accepted
);
if matches!(owner.1.case, Case::Ready) {
let (call, _) = owner
.1
.observer
.start_external_call_checked(
"app",
"module",
"service",
"route",
Some("safe-trace"),
)
.unwrap();
let event = EventContext::new(
saddle_observability::RequestIdentity::new("request").unwrap(),
saddle_observability::RouteIdentity::new("route").unwrap(),
1,
)
.unwrap();
let admission = owner.2.take().unwrap();
assert!(owner.2.take().is_none());
context.publish(saddle_core::RequestIdentityGroup::from_validated(
call.context(), "request", "route", 1,
saddle_core::ContextFact::NotEstablished,
).unwrap()).unwrap();
let (dispatch, submission) = owner.0.take().unwrap().bind_reserved_observation(
owner.1.observer.clone(),
call.context().clone(),
event,
admission,
&context.view(),
None,
);
assert_unavailable(submission);
owner.0 = Some(dispatch);
Ok(41)
} else {
std::future::pending().await
}
}
}
fn assert_unavailable(submission: saddle_observability::AdmissionCapacitySubmission) {
use saddle_observability::DiagnosticSubmission::OutputUnavailable;
assert_eq!(submission.cpu, OutputUnavailable);
assert_eq!(submission.memory, OutputUnavailable);
assert_eq!(submission.database, OutputUnavailable);
assert_eq!(submission.profuse_contract, OutputUnavailable);
}
fn factory<'a>(
owner: &'a mut Owner,
context: ReservedTaskContext,
) -> ReservedBorrowedFuture<'a, u32> {
Box::pin(body(owner, context))
}
fn returned_layout<I, R>(_: impl FnOnce(I) -> R) -> Layout {
Layout::new::<R>()
}
fn body_layout() -> Layout {
returned_layout(|(o, c): (&'static mut Owner, ReservedTaskContext)| body(o, c))
}
#[test]
fn admission_event_survives_owner_lifecycle() {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(async {
for case in [Case::Ready, Case::PrePoll, Case::Pending, Case::Withdraw] {
let shared = Arc::new(SharedRuntimeProcess {
process: Mutex::new(Some(ProfuseGwRuntimeProcess::new(
crate::request_task::reserved::tests::process(),
))),
});
let lease = ProfuseGwProcessLease {
shared: shared.clone(),
};
let startup = lease.take_database_startup_half().unwrap();
let entered = Arc::new(AtomicBool::new(false));
let input = Input {
case,
entered: entered.clone(),
observer: Observer::with_writer(
saddle_observability::ObserverConfig::default(),
std::io::sink(),
)
.unwrap(),
};
let ReservedDispatchOutcome::Ready {
root,
future,
mut ticket,
} = lease.try_reserved_dispatch(
ContextLabel::checked("app").unwrap(),
None,
body_layout(),
&[],
input,
factory,
)
else {
panic!("real permit fit")
};
let held = root
.view(saddle_core::request_context::RequestViewPhase::Reading)
.unwrap();
let before_join = held.test_available_storage();
let joined = if matches!(case, Case::Withdraw) {
ticket.withdraw(future).ok().unwrap()
} else {
let mut tasks = tokio::task::JoinSet::new();
let handle = tasks.spawn(future);
ticket.bind(handle.id()).unwrap();
if matches!(case, Case::PrePoll) {
handle.abort();
}
if matches!(case, Case::Pending) {
while !entered.load(Ordering::SeqCst) {
tokio::task::yield_now().await;
}
handle.abort();
}
if matches!(case, Case::Ready) {
while !handle.is_finished() {
tokio::task::yield_now().await;
}
}
assert!(matches!(
lease.try_admit(),
ProfuseGwCoordinatorAdmissionOutcome::CapacityRejected(_)
));
ticket
.complete(tasks.join_next().await.unwrap())
.ok()
.unwrap()
};
drop(root);
let mut recovered = joined.recover(Default::default()).ok().unwrap();
assert!(
held.test_available_storage() > before_join,
"task storage released only after matching recovery"
);
if matches!(case, Case::Ready) {
assert_eq!(recovered.result.take().unwrap().ok(), Some(41));
assert!(recovered.owner.2.is_none());
assert!(recovered.owner.0.as_ref().unwrap().observation.is_some());
} else {
assert_eq!(
recovered.owner.2.as_ref().unwrap().receipt.decision(),
ProfuseGwCapacityDecision::Accepted
);
assert!(recovered.owner.0.as_ref().unwrap().observation.is_none());
if matches!(case, Case::PrePoll | Case::Withdraw) {
assert!(!entered.load(Ordering::SeqCst));
}
assert_unavailable(held.record_admission(
recovered.owner.2.take().unwrap(),
&recovered.owner.1.observer,
None,
saddle_observability::AdmissionConstruction::Cancelled,
));
assert!(recovered.owner.2.is_none());
}
finish_profusegw_without_database(recovered.owner.0.take().unwrap(), ())
.ok()
.unwrap();
drop(recovered);
drop(held);
match lease.try_admit() {
ProfuseGwCoordinatorAdmissionOutcome::Ready(d, _) => d.cancel(),
_ => panic!("original profile retry"),
}
drop((startup, lease));
shared
.process
.lock()
.unwrap()
.take()
.unwrap()
.finish()
.unwrap();
println!("AR_LIFECYCLE {case:?} ORIGINAL_OWNER=RECOVERED ZERO=PASS");
}
println!(
"AR_LAYOUT event={} owner={} body={} demand={}",
Layout::new::<ProfuseGwAdmissionEvent>().size(),
Layout::new::<Owner>().size(),
body_layout().size(),
crate::request_task::reserved::dispatch_storage_bytes_for::<Input, u32, _>(
&factory,
body_layout(),
&[]
)
.unwrap()
);
});
}
#[test]
fn admitted_failures_restore_event_input_and_uncalled_factory() {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(async {
for construction in [false, true] {
let shared = Arc::new(SharedRuntimeProcess {
process: Mutex::new(Some(ProfuseGwRuntimeProcess::new(
crate::request_task::reserved::tests::process(),
))),
});
let lease = ProfuseGwProcessLease {
shared: shared.clone(),
};
let startup = lease.take_database_startup_half().unwrap();
let entered = Arc::new(AtomicBool::new(false));
let input = Input {
case: Case::Ready,
entered: entered.clone(),
observer: Observer::with_writer(
saddle_observability::ObserverConfig::default(),
std::io::sink(),
)
.unwrap(),
};
let backing = [(Layout::from_size_align(lease.resource_snapshot().unwrap().framework_capacity, 1).unwrap(), 1)];
crate::request_task::reserved::FAIL_AR_CONSTRUCTION
.with(|flag| flag.set(construction));
let outcome = lease.try_reserved_dispatch(
ContextLabel::checked("app").unwrap(),
None,
body_layout(),
if construction { &[] } else { &backing },
input,
factory,
);
let ReservedDispatchOutcome::Rejected {
input,
make,
reason: ReservedDispatchRejection::Admitted { admission, failure },
} = outcome
else {
panic!("admitted failure must return original event")
};
assert_eq!(
admission.receipt.decision(),
ProfuseGwCapacityDecision::Accepted
);
assert_eq!(admission.receipt.used(), 1);
assert!(!entered.load(Ordering::SeqCst));
assert!(Arc::ptr_eq(&input.entered, &entered));
assert!(matches!(
(construction, failure),
(false, ReservedDispatchConstructionFailure::Storage(_))
| (true, ReservedDispatchConstructionFailure::Context(_))
));
assert_unavailable(admission.record_unrooted(
&input.observer,
None,
&ContextLabel::checked("app").unwrap(),
saddle_core::RequestViewPhase::Admitted,
saddle_observability::AdmissionConstruction::Failed(if construction {
saddle_observability::AdmissionConstructionFailure::Context
} else {
saddle_observability::AdmissionConstructionFailure::Storage
}),
));
let ReservedDispatchOutcome::Ready {
root,
future,
ticket,
} = lease.try_reserved_dispatch(
ContextLabel::checked("app").unwrap(),
None,
body_layout(),
&[],
input,
make,
)
else {
panic!("withdrawn resources and exact factory usable")
};
let mut recovered = ticket
.withdraw(future)
.ok()
.unwrap()
.recover(Default::default())
.ok()
.unwrap();
assert_eq!(
recovered.owner.2.as_ref().unwrap().receipt.decision(),
ProfuseGwCapacityDecision::Accepted
);
recovered.owner.0.take().unwrap().cancel();
drop((recovered, root, startup, lease));
shared
.process
.lock()
.unwrap()
.take()
.unwrap()
.finish()
.unwrap();
println!(
"AR_FAILURE construction={construction} INPUT_FACTORY_EVENT=RESTORED ZERO=PASS"
);
}
});
}
type RefusalOwner = Option<ProfuseGwAdmissionEvent>;
#[test]
fn synchronous_capacity_rejection_reads_original_atomic_resource_diagnostic() {
let status = std::process::Command::new(std::env::current_exe().unwrap())
.arg("synchronous_capacity_rejection_diagnostic_child")
.arg("--ignored")
.env("SADDLE_I027_REJECTION_CHILD", "1")
.status().unwrap();
assert!(status.success(), "isolated original rejection readback");
}
#[test]
#[ignore = "run by the isolated parent test"]
fn synchronous_capacity_rejection_diagnostic_child() {
if std::env::var_os("SADDLE_I027_REJECTION_CHILD").is_none() { return; }
tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap().block_on(async {
let directory = std::env::temp_dir().join(format!("saddle-i027-rejection-{}", std::process::id()));
std::fs::create_dir_all(&directory).unwrap();
let mut writer = saddle_observability::EmergencyDiagnostics::start(
&saddle_observability::FileLoggingConfig::new(&directory, saddle_observability::Rotation::Daily)
).unwrap();
let output = writer.handle();
let observer = Observer::with_writer(Default::default(), std::io::sink()).unwrap();
let shared = Arc::new(SharedRuntimeProcess { process: Mutex::new(Some(ProfuseGwRuntimeProcess::new(
crate::request_task::reserved::tests::process()))) });
let lease = ProfuseGwProcessLease { shared: shared.clone() };
let startup = lease.take_database_startup_half().unwrap();
let ProfuseGwCoordinatorAdmissionOutcome::Ready(active, accepted) = lease.try_admit() else { panic!("first active route") };
let ProfuseGwCoordinatorAdmissionOutcome::CapacityRejected(event) = lease.try_admit() else { panic!("original synchronous capacity rejection") };
let application = ContextLabel::checked("app").unwrap();
let submissions = event.record_capacity_context(
&observer, Some(&output),
saddle_observability::AdmissionEventContext::Unrooted {
application: &application, lifecycle: saddle_core::RequestViewPhase::Admitted,
},
saddle_observability::AdmissionConstruction::NotAttempted, None,
);
assert!(matches!(submissions.cpu, saddle_observability::DiagnosticSubmission::Enqueued));
assert!(matches!(submissions.memory, saddle_observability::DiagnosticSubmission::Enqueued));
assert!(matches!(submissions.database, saddle_observability::DiagnosticSubmission::Enqueued));
assert!(matches!(submissions.profuse_contract, saddle_observability::DiagnosticSubmission::Enqueued));
active.cancel();
drop((accepted, startup, lease, output));
let until = std::time::Instant::now() + std::time::Duration::from_secs(3);
while writer.shutdown() != saddle_observability::DiagnosticShutdown::Finished {
assert!(std::time::Instant::now() < until);
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
let text = std::fs::read_to_string(directory.join("saddle.emergency.log")).unwrap();
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.capacity").collect();
assert_eq!(rows.len(), 4);
assert_eq!(rows.iter().map(|row| row["capacity_dimension"].as_str().unwrap()).collect::<Vec<_>>(),
["cpu", "memory", "database", "profuse_contract"]);
let resource = &rows[0]["resource"];
assert_eq!(resource["db_available"], false);
assert!(rows.iter().all(|row|
row["decision"] == "capacity_rejected"
&& row["outcome"] == "rejected"
&& row["reject_reason"] == "at_limit"
&& row["construction"]["status"] == "not_attempted"
&& row["resource_source"] == "admission_ledger_same_lock"
&& &row["resource"] == resource
&& row["context"]["application"]["value"] == "app"
&& row["context"]["request"]["state"] == "not_established"));
shared.process.lock().unwrap().take().unwrap().finish().unwrap();
});
}
fn refusal<'a>(
owner: &'a mut RefusalOwner,
_: ReservedTaskContext,
) -> ReservedBorrowedFuture<'a, u32> {
Box::pin(async move {
assert!(owner.is_some());
Ok(41)
})
}
#[test]
fn rejection_growth_keeps_original_event_until_last_reference() {
tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap().block_on(async {
let shared=Arc::new(SharedRuntimeProcess{process:Mutex::new(Some(ProfuseGwRuntimeProcess::new(crate::request_task::reserved::tests::process())))});
let lease=ProfuseGwProcessLease{shared:shared.clone()};
let startup=lease.take_database_startup_half().unwrap();
let ProfuseGwCoordinatorAdmissionOutcome::Ready(active,accepted)=lease.try_admit() else {panic!("active")};
let mut slots=Vec::new();
for _ in 0..33 {
let ReservedDispatchOutcome::Rejected{reason:ReservedDispatchRejection::Capacity(event),..}=lease.try_reserved_dispatch(ContextLabel::checked("app").unwrap(),None,body_layout(),&[],Input{case:Case::Ready,entered:Arc::new(AtomicBool::new(false)),observer:Observer::with_writer(saddle_observability::ObserverConfig::default(),std::io::sink()).unwrap()},factory) else {panic!("original rejected event")};
assert_eq!(event.receipt.decision(),ProfuseGwCapacityDecision::CapacityRejected);
slots.push(lease.try_reserved_rejection(ContextLabel::checked("app").unwrap(),None,Layout::new::<[u8;64]>(),&[],Some(event),refusal).ok().expect("actual storage can cross old rejection slot limit"));
}
let held=slots[0].0.view(saddle_core::request_context::RequestViewPhase::Reading).unwrap();
for (root,future,ticket) in slots {
let recovered=ticket.withdraw(future).ok().unwrap().recover(Default::default()).ok().unwrap();
assert_eq!(recovered.owner.as_ref().unwrap().receipt.decision(),ProfuseGwCapacityDecision::CapacityRejected);
drop((recovered,root));
}
let before=lease.resource_snapshot().unwrap().framework_charged;
drop(held);
let after=lease.resource_snapshot().unwrap();
assert!(after.framework_charged<before,"last root view must release its actual storage");
assert!(after.framework_charged>0,"live registry and active request stay charged");
assert!(after.healthy);
active.cancel();drop((accepted,startup,lease));shared.process.lock().unwrap().take().unwrap().finish().unwrap();
});
}