use super::*;
use std::{
pin::pin,
sync::atomic::{AtomicUsize, Ordering},
task::{Context, Poll, Waker},
};
static BUILDS: AtomicUsize = AtomicUsize::new(0);
static DROPS: AtomicUsize = AtomicUsize::new(0);
static DROP_MINIMUM: AtomicUsize = AtomicUsize::new(0);
static DROP_VIEW: std::sync::Mutex<Option<saddle_runtime::profusegw::ProfuseGwResourceView>> =
std::sync::Mutex::new(None);
pub(super) fn inner_constructed() {
BUILDS.fetch_add(1, Ordering::SeqCst);
}
crate::database_operations! {
relational;
pub mod fixtures {
namespace "control.named";
schema { "named_records" as records { "id" as id:u64, "stamp" as stamp:timestamp_utc(0) } }
unique records { id };
update Write { parameters WriteParams {id:u64,stamp:timestamp_utc(0)}
table records; values {stamp:bind(stamp)} where records.id==bind(id); }
query Rows { parameters RowsParams {id:u64,stamp:timestamp_utc(0)}
from records as r; where r.id==bind(id) && r.stamp==bind(stamp);
result rows RowsRow {id:r.id} }
query Optional { parameters OptionalParams {id:u64,stamp:timestamp_utc(0)}
from records as r; where r.id==bind(id) && r.stamp==bind(stamp);
result optional OptionalRow {id:r.id} }
query Probe { parameters ProbeParams {id:u64,stamp:timestamp_utc(0)}
from records as r; where r.id==bind(id) && r.stamp==bind(stamp);
result probe ProbeRow {id:r.id} }
query Page { parameters PageParams {}
from records as r;
order {r.id asc nulls_last} result page PageRow {id:r.id,stamp:r.stamp} }
query Closed { parameters ClosedParams {id:u64}
from records as r; where r.id==bind(id);
result optional ClosedRow {id:r.id} }
}
}
macro_rules! dropped_parameters {
($($p:ty),*) => {$ (impl Drop for $p {
fn drop(&mut self) {
let minimum = DROP_MINIMUM.load(Ordering::SeqCst);
if minimum != 0 {
assert!(DROP_VIEW.lock().unwrap().as_ref().unwrap().snapshot().unwrap().framework_charged >= minimum,
"physical parameter destruction occurs while the stored future is still charged");
}
DROPS.fetch_add(1, Ordering::SeqCst);
}
})*};
}
dropped_parameters!(
fixtures::WriteParams,
fixtures::RowsParams,
fixtures::OptionalParams,
fixtures::ProbeParams
);
static PAGE_PHASE: AtomicUsize = AtomicUsize::new(0);
static PAGE_CONVERSION_DROPS: AtomicUsize = AtomicUsize::new(0);
static PAGE_ORIGINAL_DROPS: AtomicUsize = AtomicUsize::new(0);
impl Drop for fixtures::PageParams {
fn drop(&mut self) {
let phase = PAGE_PHASE.load(Ordering::SeqCst);
if phase != 0 {
let charged = DROP_VIEW
.lock()
.unwrap()
.as_ref()
.unwrap()
.snapshot()
.unwrap()
.framework_charged;
assert!(charged >= DROP_MINIMUM.load(Ordering::SeqCst));
match phase {
1 => {
PAGE_CONVERSION_DROPS.fetch_add(1, Ordering::SeqCst);
}
2 => {
PAGE_ORIGINAL_DROPS.fetch_add(1, Ordering::SeqCst);
}
_ => panic!("unexpected page lifecycle phase"),
}
}
DROPS.fetch_add(1, Ordering::SeqCst);
}
}
fn stamp(year: u16) -> saddle_db::internal::DbTimestamp<0> {
saddle_db::internal::DbTimestamp::new(
saddle_db::internal::DbDate::new(year, 1, 1).unwrap(),
0,
0,
1,
0,
)
.unwrap()
}
fn page(_year: u16) -> crate::database_rows::PageRequest<fixtures::Page> {
crate::database_rows::PageRequest::first(fixtures::PageParams {}, 1).unwrap()
}
fn send<F: Future + Send>(_: &F) {}
impl DatabaseRequest {
pub(crate) async fn control_named_storage_admission(
&mut self,
view: &saddle_runtime::profusegw::ProfuseGwResourceView,
surface: usize,
) {
assert!(self.state.reserved_context.is_some() && self.process.is_some());
let before = view.snapshot().unwrap().framework_charged;
let start_drops = DROPS.load(Ordering::SeqCst);
macro_rules! unpolled {
($call:expr) => {{
let f = $call;
send(&f);
drop(f);
}};
}
unpolled!(self.named_write::<fixtures::Write>(fixtures::WriteParams {
id: 7,
stamp: stamp(2020)
}));
unpolled!(self.named_rows::<fixtures::Rows>(fixtures::RowsParams {
id: 7,
stamp: stamp(2020)
}));
unpolled!(
self.named_optional::<fixtures::Optional>(fixtures::OptionalParams {
id: 7,
stamp: stamp(2020)
})
);
unpolled!(self.named_probe::<fixtures::Probe>(fixtures::ProbeParams {
id: 7,
stamp: stamp(2020)
}));
unpolled!(self.named_page::<fixtures::Page>(page(2020)));
assert_eq!(DROPS.load(Ordering::SeqCst) - start_drops, 5);
assert_eq!(BUILDS.load(Ordering::SeqCst), 0);
assert_eq!(view.snapshot().unwrap().framework_charged, before);
macro_rules! reject {
($call:expr) => {{
let result = $call.await;
let Err(QueryFailure::Target { parameters, error }) = result else {
panic!("original target rejection required")
};
assert_eq!(error, saddle_db::internal::TimestampTargetError::OutOfRange);
assert_eq!(parameters.id, 7);
assert_eq!(parameters.stamp.as_sql_str(), stamp(2107).as_sql_str());
drop(parameters);
}};
}
reject!(self.named_write::<fixtures::Write>(fixtures::WriteParams {
id: 7,
stamp: stamp(2107)
}));
reject!(self.named_rows::<fixtures::Rows>(fixtures::RowsParams {
id: 7,
stamp: stamp(2107)
}));
reject!(
self.named_optional::<fixtures::Optional>(fixtures::OptionalParams {
id: 7,
stamp: stamp(2107)
})
);
reject!(self.named_probe::<fixtures::Probe>(fixtures::ProbeParams {
id: 7,
stamp: stamp(2107)
}));
assert_eq!(DROPS.load(Ordering::SeqCst) - start_drops, 9);
assert_eq!(BUILDS.load(Ordering::SeqCst), 0);
assert_eq!(view.snapshot().unwrap().framework_charged, before);
assert!(self.state.reserved_failure.lock().unwrap().is_none());
assert!(matches!(
*self.state.terminal.lock().unwrap(),
RequestTerminal::Untouched(_)
));
let original_process = self.process.take();
let private =
self.named_optional_inner::<fixtures::Closed>(fixtures::ClosedParams { id: 7 });
let bytes = std::mem::size_of_val(&private);
drop(private);
{
let mut future =
pin!(self.named_optional::<fixtures::Closed>(fixtures::ClosedParams { id: 7 }));
assert!(matches!(
future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop())),
Poll::Ready(Err(QueryFailure::Database(ScopeDatabaseError::State)))
));
assert_eq!(view.snapshot().unwrap().framework_charged - before, bytes);
drop(future); }
assert_eq!(view.snapshot().unwrap().framework_charged, before);
self.process = original_process;
struct Held<const N: usize>([u8; N]);
impl<const N: usize> Future for Held<N> {
type Output = ();
fn poll(self: std::pin::Pin<&mut Self>, _: &mut Context<'_>) -> Poll<()> {
std::hint::black_box(&self.0);
Poll::Pending
}
}
let memory = self.framework_future_memory();
macro_rules! fill {
($n:expr) => {{
let mut held = Vec::new();
loop {
match saddle_admission::managed_framework_future(
&memory,
|| Held::<$n>([0; $n]),
|_| {},
) {
Ok(f) => held.push(f),
Err(saddle_admission::AdmissionError::FrameworkReserveExceeded {
..
}) => break,
Err(e) => panic!("unexpected reserve failure {e:?}"),
}
assert!(held.len() < 1024);
}
held
}};
}
let large = fill!(32768);
let medium = fill!(1024);
let small = fill!(64);
let last = fill!(1);
macro_rules! denied {
($call:expr) => {{
assert!(matches!(
$call.await,
Err(QueryFailure::Database(ScopeDatabaseError::Resource))
));
}};
}
let builds = BUILDS.load(Ordering::SeqCst);
match surface {
0 => denied!(self.named_write::<fixtures::Write>(fixtures::WriteParams {
id: 7,
stamp: stamp(2020)
})),
1 => denied!(self.named_rows::<fixtures::Rows>(fixtures::RowsParams {
id: 7,
stamp: stamp(2020)
})),
2 => denied!(
self.named_optional::<fixtures::Optional>(fixtures::OptionalParams {
id: 7,
stamp: stamp(2020)
})
),
3 => denied!(self.named_probe::<fixtures::Probe>(fixtures::ProbeParams {
id: 7,
stamp: stamp(2020)
})),
4 => denied!(self.named_page::<fixtures::Page>(page(2020))),
_ => panic!("unknown surface"),
}
assert_eq!(
BUILDS.load(Ordering::SeqCst),
builds,
"no inner construction before reserve"
);
assert!(matches!(
*self.state.terminal.lock().unwrap(),
RequestTerminal::Untouched(_)
));
let occurrence = {
let slot = self.state.reserved_failure.lock().unwrap();
let e = slot.as_ref().unwrap();
assert!(e.original_was_written());
serde_json::to_value(e.occurrence()).unwrap()
};
drop((large, medium, small, last));
macro_rules! blocked {
($call:expr) => {
assert!(matches!(
$call.await,
Err(QueryFailure::Database(ScopeDatabaseError::State))
));
};
}
blocked!(self.named_write::<fixtures::Write>(fixtures::WriteParams {
id: 7,
stamp: stamp(2020)
}));
blocked!(self.named_rows::<fixtures::Rows>(fixtures::RowsParams {
id: 7,
stamp: stamp(2020)
}));
blocked!(
self.named_optional::<fixtures::Optional>(fixtures::OptionalParams {
id: 7,
stamp: stamp(2020)
})
);
blocked!(self.named_probe::<fixtures::Probe>(fixtures::ProbeParams {
id: 7,
stamp: stamp(2020)
}));
blocked!(self.named_page::<fixtures::Page>(page(2020)));
assert_eq!(BUILDS.load(Ordering::SeqCst), builds);
assert_eq!(
serde_json::to_value(
self.state
.reserved_failure
.lock()
.unwrap()
.as_ref()
.unwrap()
.occurrence()
)
.unwrap(),
occurrence
);
assert!(matches!(
*self.state.terminal.lock().unwrap(),
RequestTerminal::Untouched(_)
));
println!(
"NAMED_STORAGE_ADMISSION surface={surface} PASS unpolled/target-intact/no-inner/noSQL/RequestDb/original-preserved"
);
}
}
impl DatabaseRequest {
pub(crate) fn control_named_storage_config(
config: saddle_db::DatabaseConfig,
) -> saddle_db::DatabaseConfig {
config
.register_owned_write::<fixtures::Write>()
.register_owned_query::<fixtures::Rows>()
.register_owned_query::<fixtures::Optional>()
.register_owned_query::<fixtures::Probe>()
.register_owned_query::<fixtures::Page>()
}
pub(crate) async fn control_named_storage_normal(&mut self) {
let write = self
.named_write::<fixtures::Write>(fixtures::WriteParams {
id: 7,
stamp: stamp(2020),
})
.await;
if let Err(QueryFailure::Database(error)) = &write {
panic!("named normal write database error: {error:?}");
}
assert!(write.is_ok());
drop(write);
let rows = self
.named_rows::<fixtures::Rows>(fixtures::RowsParams {
id: 7,
stamp: stamp(2020),
})
.await;
assert!(rows.is_ok());
drop(rows);
let row = self
.named_optional::<fixtures::Optional>(fixtures::OptionalParams {
id: 7,
stamp: stamp(2020),
})
.await;
assert!(matches!(row, Ok(Some(_))));
drop(row);
let probe = self
.named_probe::<fixtures::Probe>(fixtures::ProbeParams {
id: 7,
stamp: stamp(2020),
})
.await;
assert!(matches!(
probe,
Ok(crate::database::MatchCardinality::One(_))
));
drop(probe);
let page = self.named_page::<fixtures::Page>(page(2020)).await;
assert!(page.is_ok());
drop(page);
assert!(self.state.reserved_failure.lock().unwrap().is_none());
println!("NAMED_STORAGE_NORMAL PASS five actual SQL surfaces");
}
}
impl DatabaseRequest {
pub(crate) async fn control_named_storage_database_failure(&mut self) {
let result = self
.named_write::<fixtures::Write>(fixtures::WriteParams {
id: 7,
stamp: stamp(2020),
})
.await;
assert!(matches!(
result,
Err(QueryFailure::Database(ScopeDatabaseError::Operation))
));
assert!(self.state.reserved_failure.lock().unwrap().is_none());
println!(
"NAMED_STORAGE_DATABASE_FAILURE PASS Operation/execution-owner/admission-slot-untouched"
);
}
pub(crate) async fn control_named_storage_cancel(
&mut self,
view: &saddle_runtime::profusegw::ProfuseGwResourceView,
) {
let before = view.snapshot().unwrap().framework_charged;
type RetainedShape = Option<std::sync::Arc<()>>;
let prefix = std::alloc::Layout::new::<[std::sync::atomic::AtomicUsize; 2]>()
.extend(std::alloc::Layout::new::<(
saddle_core::request_context::RequestExecutionView,
Option<RetainedShape>,
Option<RetainedShape>,
RetainedShape,
)>())
.unwrap()
.0;
let [_, (_, payload), _] = saddle_core::request_context::request_context_layouts();
let held_view = saddle_admission::StorageDemand::embedded(
prefix,
std::alloc::Layout::new::<()>(),
&[(payload, 1)],
)
.unwrap()
.bytes();
let held_scope = saddle_runtime::request_task::reserved::supervised_scope_layouts()
.into_iter()
.find(|(name, _)| *name == "scope_shared_allocation")
.unwrap()
.1
.size();
let retained = held_view + held_scope;
let private = self.named_page_inner::<fixtures::Page>(page(2020));
let bytes = std::mem::size_of_val(&private);
drop(private);
let state = self.state.clone();
let drops = DROPS.load(Ordering::SeqCst);
*DROP_VIEW.lock().unwrap() = Some(view.clone());
DROP_MINIMUM.store(before + bytes, Ordering::SeqCst);
PAGE_CONVERSION_DROPS.store(0, Ordering::SeqCst);
PAGE_ORIGINAL_DROPS.store(0, Ordering::SeqCst);
PAGE_PHASE.store(1, Ordering::SeqCst);
{
let mut future = pin!(self.named_page::<fixtures::Page>(page(2020)));
assert!(
future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop()))
.is_pending()
);
assert!(
!matches!(
*state.terminal.lock().unwrap(),
RequestTerminal::Untouched(_)
),
"real reserved serial owner entered before cancellation"
);
assert!(view.snapshot().unwrap().framework_charged >= before + bytes);
assert_eq!(PAGE_CONVERSION_DROPS.load(Ordering::SeqCst), 1);
assert_eq!(PAGE_ORIGINAL_DROPS.load(Ordering::SeqCst), 0);
PAGE_PHASE.store(2, Ordering::SeqCst);
}
PAGE_PHASE.store(0, Ordering::SeqCst);
assert_eq!(PAGE_CONVERSION_DROPS.load(Ordering::SeqCst), 1);
assert_eq!(PAGE_ORIGINAL_DROPS.load(Ordering::SeqCst), 1);
DROP_MINIMUM.store(0, Ordering::SeqCst);
*DROP_VIEW.lock().unwrap() = None;
assert_eq!(DROPS.load(Ordering::SeqCst) - drops, 2);
assert!(
matches!(
*state.terminal.lock().unwrap(),
RequestTerminal::PostDatabase(_)
),
"original ReservedDbOwner restored physical terminal"
);
assert_eq!(
view.snapshot().unwrap().framework_charged,
before + retained,
"named F fully refunded; original scope view/observation remain physically owned"
);
assert!(view.snapshot().unwrap().healthy);
println!(
"NAMED_STORAGE_CANCEL PASS actual-owner/Pending/physical-parameter-before-refund/original-terminal"
);
}
}
impl DatabaseRequest {
pub(crate) async fn control_named_storage_legacy(&mut self) {
assert!(self.state.reserved_context.is_none());
assert!(self.state.legacy.is_some());
self.control_named_storage_normal().await;
println!(
"NAMED_STORAGE_LEGACY PASS original nonreserved entry/tasks/normal-five/charged-memory"
);
}
}