use super::*;
pub(super) enum NextScope {
Continuation(saddle_core::DbScopeContinuation),
Pair(DbPhysicalRequestHalf, DbPhysicalExecutionHalf),
}
#[doc(hidden)]
pub struct ProfuseGwSerialScope<C> {
lease: ProfuseGwConcreteDbRequestLease,
cancel: Pin<Box<C>>,
entry: ScopeEntry,
}
enum ScopeEntry {
Entered,
ResumedUnused,
}
#[doc(hidden)]
pub type ProfuseGwTransactionObservation =
saddle_core::DbScopeObservation<(Observer, CallContext, EventContext)>;
#[derive(Debug, PartialEq, Eq)]
#[doc(hidden)]
pub enum ProfuseGwTransactionObservationError {
MissingContext,
Unavailable,
}
impl<C> ProfuseGwSerialScope<C> {}
impl ProfuseGwConcreteDbRequestLease {
#[doc(hidden)]
pub fn into_serial_scope<C: Future<Output = ()>>(self, cancel: C) -> ProfuseGwSerialScope<C> {
ProfuseGwSerialScope {
lease: self,
cancel: Box::pin(cancel),
entry: ScopeEntry::Entered,
}
}
}
struct Control<'a, C> {
deadline: &'a mut ProfuseGwManagedDeadline,
stop: &'a mut Option<ProfuseGwScopeStop>,
cancel: Pin<&'a mut C>,
supervising: bool,
}
impl<C: Future<Output = ()>> Control<'_, C> {
fn check(&mut self, cx: &mut Context<'_>) -> Result<(), ProfuseGwScopeStop> {
if let Some(stop) = *self.stop {
return Err(stop);
}
let stop = if self.cancel.as_mut().poll(cx).is_ready() {
Some(ProfuseGwScopeStop::Cancelled)
} else if tokio::time::Instant::now() >= self.deadline.timer.deadline()
|| self.deadline.timer.as_mut().poll(cx).is_ready()
{
Some(ProfuseGwScopeStop::TimedOut)
} else {
None
};
*self.stop = stop;
stop.map_or(Ok(()), Err)
}
}
#[doc(hidden)]
pub struct ProfuseGwScopeDriver<'a, C> {
execution: &'a ProfuseGwLightweightExecutionOwner,
control: Mutex<Control<'a, C>>,
}
#[derive(Debug, PartialEq, Eq)]
#[doc(hidden)]
pub enum ProfuseGwScopeFailure {
Stopped(ProfuseGwScopeStop),
Panicked,
AlreadySupervised,
ScopeAlreadyEntered,
}
impl<C: Future<Output = ()>> ProfuseGwSerialScope<C> {
#[doc(hidden)]
pub fn take_transaction_observation(
&mut self,
) -> Result<ProfuseGwTransactionObservation, ProfuseGwTransactionObservationError> {
let observation = self
.lease
.observation
.as_ref()
.ok_or(ProfuseGwTransactionObservationError::MissingContext)?;
let context = (
observation.observer.clone(),
observation.context.clone(),
observation.event_context.clone(),
);
let correlation = self
.lease
.db_request
.take_scope_observation(&self.lease.db_execution, context)
.map_err(|_| ProfuseGwTransactionObservationError::Unavailable)?;
self.entry = ScopeEntry::Entered;
Ok(correlation)
}
#[doc(hidden)]
pub fn driver(&mut self) -> ProfuseGwScopeDriver<'_, C> {
self.entry = ScopeEntry::Entered;
self.borrow_driver()
}
fn borrow_driver(&mut self) -> ProfuseGwScopeDriver<'_, C> {
ProfuseGwScopeDriver {
execution: &self.lease.execution,
control: Mutex::new(Control {
deadline: &mut self.lease.deadline,
stop: &mut self.lease.scope_stop,
cancel: self.cancel.as_mut(),
supervising: false,
}),
}
}
#[doc(hidden)]
pub async fn supervise_between<F: Future>(
&mut self,
future: F,
) -> Result<F::Output, ProfuseGwScopeFailure> {
if matches!(self.entry, ScopeEntry::Entered) {
return Err(ProfuseGwScopeFailure::ScopeAlreadyEntered);
}
self.borrow_driver().supervise(future).await
}
#[doc(hidden)]
#[allow(
clippy::result_large_err,
reason = "return the exact linear whole and value without a new allocation"
)]
pub fn finish_unentered_response<T>(self, value: T) -> Result<T, (Self, T)> {
if matches!(self.entry, ScopeEntry::Entered) {
return Err((self, value));
}
let ProfuseGwConcreteDbRequestLease {
execution,
deadline,
db_request,
db_execution,
observation,
scope_stop,
} = self.lease;
match seal_db_request_not_used(db_request, db_execution, value) {
Ok(receipt) => {
record_terminal(observation, execution.cancel_observed());
drop(deadline);
Ok(receipt.into_value())
}
Err((db_request, db_execution, value)) => Err((
Self {
lease: ProfuseGwConcreteDbRequestLease {
execution,
deadline,
db_request,
db_execution,
observation,
scope_stop,
},
cancel: self.cancel,
entry: self.entry,
},
value,
)),
}
}
#[doc(hidden)]
pub fn into_physical_finalization(
self,
) -> (ProfuseGwSerialCompletion<C>, DbPhysicalExecutionHalf) {
let (completion, execution) = self
.lease
.into_physical_finalization()
.into_database_execution();
(
ProfuseGwSerialCompletion {
completion,
cancel: self.cancel,
},
execution,
)
}
}
impl<C: Future<Output = ()>> ProfuseGwScopeDriver<'_, C> {
#[doc(hidden)]
pub async fn checkpoint(&self) -> Result<(), ProfuseGwScopeStop> {
let mut yielded = false;
poll_fn(|cx| {
let mut control = self.control.lock().unwrap_or_else(|p| p.into_inner());
if let Err(stop) = control.check(cx) {
return Poll::Ready(Err(stop));
}
if !yielded {
yielded = true;
cx.waker().wake_by_ref();
Poll::Pending
} else {
Poll::Ready(Ok(()))
}
})
.await
}
#[doc(hidden)]
pub async fn supervise<F: Future>(
&self,
future: F,
) -> Result<F::Output, ProfuseGwScopeFailure> {
{
let mut control = self.control.lock().unwrap_or_else(|p| p.into_inner());
if control.supervising {
return Err(ProfuseGwScopeFailure::AlreadySupervised);
}
control.supervising = true;
}
let mut guard = SupervisionGuard {
control: &self.control,
completed: false,
};
let mut future = Box::pin(future);
let outcome = poll_fn(|cx| {
{
let mut control = self.control.lock().unwrap_or_else(|p| p.into_inner());
if let Err(stop) = control.check(cx) {
return Poll::Ready(Err(ProfuseGwScopeFailure::Stopped(stop)));
}
}
let mut caught =
poll_fn(
|cx| match catch_unwind(AssertUnwindSafe(|| future.as_mut().poll(cx))) {
Ok(Poll::Pending) => Poll::Pending,
Ok(Poll::Ready(value)) => Poll::Ready(Ok(value)),
Err(_) => Poll::Ready(Err(())),
},
);
match self
.execution
.poll_database_query(Pin::new(&mut caught), cx)
{
Poll::Pending => Poll::Pending,
Poll::Ready(Ok(value)) => {
let mut control = self.control.lock().unwrap_or_else(|p| p.into_inner());
let _ = control.check(cx);
Poll::Ready(Ok(value))
}
Poll::Ready(Err(())) => {
let mut control = self.control.lock().unwrap_or_else(|p| p.into_inner());
control.stop.get_or_insert(ProfuseGwScopeStop::Cancelled);
Poll::Ready(Err(ProfuseGwScopeFailure::Panicked))
}
}
})
.await;
drop(future);
guard.completed = true;
outcome
}
}
struct SupervisionGuard<'a, 'b, C> {
control: &'a Mutex<Control<'b, C>>,
completed: bool,
}
impl<C> Drop for SupervisionGuard<'_, '_, C> {
fn drop(&mut self) {
let mut control = self.control.lock().unwrap_or_else(|p| p.into_inner());
if !self.completed {
control.stop.get_or_insert(ProfuseGwScopeStop::Cancelled);
}
control.supervising = false;
}
}
#[doc(hidden)]
pub struct ProfuseGwSerialCompletion<C> {
completion: ProfuseGwDatabaseFinalizationCompletion,
cancel: Pin<Box<C>>,
}
impl<C: Future<Output = ()>> ProfuseGwSerialCompletion<C> {
#[doc(hidden)]
pub fn poll_physical_stop(&mut self, cx: &mut Context<'_>) -> Poll<()> {
let mut control = Control {
deadline: &mut self.completion.deadline,
stop: &mut self.completion.scope_stop,
cancel: self.cancel.as_mut(),
supervising: false,
};
if control.check(cx).is_err() {
Poll::Ready(())
} else {
Poll::Pending
}
}
#[doc(hidden)]
#[allow(
clippy::result_large_err,
reason = "foreign receipt must return both complete owners"
)]
pub fn complete<T>(
self,
physical: DbPhysicalDispositionOwner<T>,
) -> Result<ProfuseGwSuspendedScope<C, T>, (Self, DbPhysicalDispositionOwner<T>)> {
match finish_profusegw_database_disposition(self.completion, physical) {
Ok(owner) => Ok(ProfuseGwSuspendedScope {
inner: Some((owner, self.cancel)),
}),
Err(failure) => Err((
Self {
completion: failure.completion,
cancel: self.cancel,
},
failure.physical,
)),
}
}
}
#[doc(hidden)]
pub struct ProfuseGwSuspendedScope<C, T> {
inner: Option<(ProfuseGwPostDatabaseManagedOwner<T>, Pin<Box<C>>)>,
}
#[derive(Debug)]
#[doc(hidden)]
pub enum ProfuseGwScopeResumeError {
Consumed,
Stopped(ProfuseGwScopeStop),
ScopeExhausted,
Admission(AdmissionError),
}
impl<C: Future<Output = ()>, T> ProfuseGwSuspendedScope<C, T> {
#[doc(hidden)]
pub async fn resume(
&mut self,
) -> Result<(ProfuseGwSerialScope<C>, T), ProfuseGwScopeResumeError> {
let (owner, cancel) = self
.inner
.as_mut()
.ok_or(ProfuseGwScopeResumeError::Consumed)?;
scope_checkpoint(
&mut owner.terminal.deadline,
&mut owner.terminal.scope_stop,
cancel.as_mut(),
)
.await
.map_err(ProfuseGwScopeResumeError::Stopped)?;
let Some((owner, cancel)) = self.inner.take() else {
return Err(ProfuseGwScopeResumeError::Consumed);
};
let ProfuseGwPostDatabaseManagedOwner { terminal, value } = owner;
let ProfuseGwPostDatabaseRequestTerminal {
admission,
deadline,
observation,
scope_stop,
next_scope,
} = terminal;
let pair = match next_scope {
NextScope::Pair(request, execution) => (request, execution),
NextScope::Continuation(continuation) => match continuation.into_next_scope() {
Ok(pair) => pair,
Err(continuation) => {
self.inner = Some((
ProfuseGwPostDatabaseManagedOwner {
terminal: ProfuseGwPostDatabaseRequestTerminal {
admission,
deadline,
observation,
scope_stop,
next_scope: NextScope::Continuation(continuation),
},
value,
},
cancel,
));
return Err(ProfuseGwScopeResumeError::ScopeExhausted);
}
},
};
match admission.resume_after_runtime_checkpoint() {
Ok(execution) => Ok((
ProfuseGwSerialScope {
lease: ProfuseGwConcreteDbRequestLease {
execution,
deadline,
db_request: pair.0,
db_execution: pair.1,
observation,
scope_stop,
},
cancel,
entry: ScopeEntry::ResumedUnused,
},
value,
)),
Err((admission, error)) => {
self.inner = Some((
ProfuseGwPostDatabaseManagedOwner {
terminal: ProfuseGwPostDatabaseRequestTerminal {
admission,
deadline,
observation,
scope_stop,
next_scope: NextScope::Pair(pair.0, pair.1),
},
value,
},
cancel,
));
Err(ProfuseGwScopeResumeError::Admission(error))
}
}
}
#[doc(hidden)]
#[allow(
clippy::result_large_err,
reason = "preserve the inert whole on repeated consumption"
)]
pub fn into_response_parts(
mut self,
) -> Result<(T, ProfuseGwPostDatabaseRequestTerminal), Self> {
match self.inner.take() {
Some((owner, _cancel)) => Ok(owner.into_response_parts()),
None => Err(self),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicBool, Ordering};
fn process() -> ProfuseGwRuntimeProcess {
process_capacity(1)
}
fn process_capacity(capacity: u32) -> ProfuseGwRuntimeProcess {
let pending = saddle_admission::freeze_deployment_resource_budget(
capacity, 32, 5_000, capacity, capacity, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
)
.unwrap();
let (application, 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(5_000),
))
.ok()
.unwrap();
let (whole, receipt) = saddle_core::pair_bootstrap_rendezvous(application, listener)
.ok()
.unwrap();
let budget =
saddle_admission::bind_deployment_resource_budget_bootstrap(pending, whole, receipt)
.ok()
.unwrap();
ProfuseGwRuntimeProcess::new(prepare_profusegw_lightweight_profile(budget).ok().unwrap())
}
fn admit(process: &ProfuseGwRuntimeProcess) -> ProfuseGwConcreteDbRequestLease {
match process.try_admit() {
ProfuseGwCoordinatorAdmissionOutcome::Ready(dispatch, _) => {
dispatch.into_database_request()
}
_ => panic!("admission must succeed"),
}
}
fn full(process: &ProfuseGwRuntimeProcess) {
assert!(matches!(
process.try_admit(),
ProfuseGwCoordinatorAdmissionOutcome::CapacityRejected(_)
));
}
#[tokio::test]
async fn scope_observation_original_context_serial_and_zero() {
let mut process = process();
let physical = process.db_startup.take().unwrap().into_process_capability();
let mut lease = admit(&process);
let original_deadline = lease.deadline.timer.deadline();
let observer = Observer::with_writer(
saddle_observability::ObserverConfig::default(),
std::io::sink(),
)
.unwrap();
let (call, _) = observer
.start_external_call_checked("app", "module", "service", "route", Some("safe-trace"))
.unwrap();
let context = call.context().clone();
let event = EventContext::new(
saddle_observability::RequestIdentity::new("request").unwrap(),
saddle_observability::RouteIdentity::new("route").unwrap(),
1,
)
.unwrap();
lease.observation = Some(ProfuseGwRequestObservation {
observer: observer.clone(),
context: context.clone(),
event_context: event.clone(),
});
let mut scope = lease.into_serial_scope(std::future::pending());
for expected in [0, 1] {
let saved = scope.lease.observation.take().unwrap();
assert!(matches!(
scope.take_transaction_observation(),
Err(ProfuseGwTransactionObservationError::MissingContext)
));
scope.lease.observation = Some(saved);
let token = scope.take_transaction_observation().unwrap();
assert!(matches!(
scope.take_transaction_observation(),
Err(ProfuseGwTransactionObservationError::Unavailable)
));
assert_eq!(scope.lease.deadline.timer.deadline(), original_deadline);
let (completion, execution) = scope.into_physical_finalization();
let checked = token.bind_terminal(&execution).ok().unwrap();
let ((_observer, actual_context, actual_event), fields) = checked.into_log_parts();
assert_eq!(actual_context, context);
assert_eq!(actual_event, event);
assert_eq!(
serde_json::to_value(fields).unwrap(),
serde_json::json!({"transaction_scope":expected})
);
let proof = physical.connection_returned(execution, ()).ok().unwrap();
let mut suspended = completion.complete(proof).ok().unwrap();
full(&process);
scope = suspended.resume().await.unwrap().0;
}
scope.finish_unentered_response(()).ok().unwrap();
process.finish().unwrap();
call.succeed();
observer.flush().await.unwrap();
}
struct Cancel {
requested: Arc<AtomicBool>,
completed: bool,
}
impl Future for Cancel {
type Output = ();
fn poll(mut self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<()> {
assert!(
!self.completed,
"completed cancellation future polled twice"
);
if self.requested.load(Ordering::SeqCst) {
self.completed = true;
Poll::Ready(())
} else {
Poll::Pending
}
}
}
struct BorrowedPending<'a>(&'a mut bool);
impl Future for BorrowedPending<'_> {
type Output = ();
fn poll(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<()> {
Poll::Pending
}
}
impl Drop for BorrowedPending<'_> {
fn drop(&mut self) {
*self.0 = true;
}
}
struct BorrowingSession<'driver, 'request, C> {
driver: &'driver ProfuseGwScopeDriver<'request, C>,
dropped: bool,
}
impl<C: Future<Output = ()> + Send> BorrowingSession<'_, '_, C> {
async fn run_body<T>(
&mut self,
body: impl for<'tx> FnOnce(&'tx mut Self) -> Pin<Box<dyn Future<Output = T> + Send + 'tx>>,
) -> T {
body(self).await
}
}
#[tokio::test]
async fn serial_scope_same_request_keeps_credit_timer_and_value() {
let mut process = process();
let physical = process.db_startup.take().unwrap().into_process_capability();
let lease = admit(&process);
let deadline = lease.deadline.timer.deadline();
let mut scope = lease.into_serial_scope(std::future::pending());
for value in [11, 22, 33] {
{
let driver = scope.driver();
let body = async {
driver.checkpoint().await.unwrap();
driver.checkpoint().await.unwrap();
value
};
assert_eq!(driver.supervise(body).await.unwrap(), value);
}
let (completion, execution) = scope.into_physical_finalization();
let disposition = physical.connection_returned(execution, value).ok().unwrap();
let mut suspended = completion.complete(disposition).ok().unwrap();
full(&process);
let (next, held_value) = suspended.resume().await.unwrap();
assert_eq!(held_value, value);
assert_eq!(next.lease.deadline.timer.deadline(), deadline);
assert!(matches!(
suspended.resume().await,
Err(ProfuseGwScopeResumeError::Consumed)
));
scope = next;
}
let (completion, execution) = scope.into_physical_finalization();
let disposition = physical
.connection_discarded(execution, "Unknown preserved")
.ok()
.unwrap();
let suspended = completion.complete(disposition).ok().unwrap();
let (value, terminal) = suspended.into_response_parts().ok().unwrap();
assert_eq!(value, "Unknown preserved");
full(&process);
finish_profusegw_after_database(terminal);
process.finish().unwrap();
}
#[tokio::test]
async fn serial_scope_non_db_pending_stop_drops_before_finalization() {
for cancelled in [false, true] {
let mut process = process();
let physical = process.db_startup.take().unwrap().into_process_capability();
let lease = admit(&process);
let requested = Arc::new(AtomicBool::new(false));
let mut scope = lease.into_serial_scope(Cancel {
requested: requested.clone(),
completed: false,
});
let dropped;
{
let driver = scope.driver();
let mut session = BorrowingSession {
driver: &driver,
dropped: false,
};
let supervised = driver.supervise(session.run_body(|session| {
Box::pin(async move {
session.driver.checkpoint().await.unwrap();
BorrowedPending(&mut session.dropped).await;
})
}));
fn require_send<T: Send>(_: &T) {}
require_send(&supervised);
{
tokio::pin!(supervised);
poll_fn(|cx| {
assert!(supervised.as_mut().poll(cx).is_pending());
Poll::Ready(())
})
.await;
poll_fn(|cx| {
assert!(supervised.as_mut().poll(cx).is_pending());
Poll::Ready(())
})
.await;
if cancelled {
requested.store(true, Ordering::SeqCst);
} else {
driver
.control
.lock()
.unwrap()
.deadline
.timer
.as_mut()
.reset(tokio::time::Instant::now());
}
assert!(matches!(
supervised.await,
Err(ProfuseGwScopeFailure::Stopped(_))
));
}
dropped = session.dropped;
}
assert!(
dropped,
"borrowed business future must die before phase finalizer"
);
let (mut completion, execution) = scope.into_physical_finalization();
poll_fn(|cx| {
assert!(completion.poll_physical_stop(cx).is_ready());
Poll::Ready(())
})
.await;
let disposition = physical.connection_discarded(execution, ()).ok().unwrap();
let mut suspended = completion.complete(disposition).ok().unwrap();
assert!(matches!(
suspended.resume().await,
Err(ProfuseGwScopeResumeError::Stopped(_))
));
assert!(matches!(
suspended.resume().await,
Err(ProfuseGwScopeResumeError::Stopped(_))
));
let (_, terminal) = suspended.into_response_parts().ok().unwrap();
finish_profusegw_after_database(terminal);
process.finish().unwrap();
}
}
#[tokio::test]
async fn serial_scope_dropped_supervision_is_sticky() {
let mut process = process();
let physical = process.db_startup.take().unwrap().into_process_capability();
let mut scope = admit(&process).into_serial_scope(std::future::pending());
let mut dropped = false;
{
let driver = scope.driver();
let mut future = Box::pin(driver.supervise(BorrowedPending(&mut dropped)));
poll_fn(|cx| {
assert!(future.as_mut().poll(cx).is_pending());
Poll::Ready(())
})
.await;
drop(future);
assert_eq!(
driver.checkpoint().await,
Err(ProfuseGwScopeStop::Cancelled)
);
}
assert!(dropped);
let (completion, execution) = scope.into_physical_finalization();
let disposition = physical.connection_discarded(execution, ()).ok().unwrap();
let mut suspended = completion.complete(disposition).ok().unwrap();
assert!(matches!(
suspended.resume().await,
Err(ProfuseGwScopeResumeError::Stopped(_))
));
let (_, terminal) = suspended.into_response_parts().ok().unwrap();
finish_profusegw_after_database(terminal);
process.finish().unwrap();
}
#[tokio::test]
async fn serial_scope_resume_pending_drop_and_foreign_return_owners() {
let mut a = process_capacity(2);
let pa = a.db_startup.take().unwrap().into_process_capability();
let (ca, ea) = admit(&a)
.into_serial_scope(std::future::pending())
.into_physical_finalization();
let (cb, eb) = admit(&a)
.into_serial_scope(std::future::pending())
.into_physical_finalization();
let ra = pa.connection_returned(ea, 1).ok().unwrap();
let rb = pa.connection_discarded(eb, 2).ok().unwrap();
let (ca, rb) = ca.complete(rb).err().unwrap();
let (cb, ra) = cb.complete(ra).err().unwrap();
full(&a);
let mut suspended = ca.complete(ra).ok().unwrap();
{
let mut attempt = Box::pin(suspended.resume());
poll_fn(|cx| {
assert!(attempt.as_mut().poll(cx).is_pending());
Poll::Ready(())
})
.await;
}
full(&a);
let (mut scope, value) = suspended.resume().await.unwrap();
assert_eq!(value, 1);
assert_eq!(scope.supervise_between(async { 7 }).await.unwrap(), 7);
full(&a);
assert_eq!(scope.finish_unentered_response(value).ok().unwrap(), 1);
let (_, terminal) = cb
.complete(rb)
.ok()
.unwrap()
.into_response_parts()
.ok()
.unwrap();
finish_profusegw_after_database(terminal);
a.finish().unwrap();
}
#[tokio::test]
async fn serial_scope_ready_loop_and_panic_keep_terminal() {
for panic_body in [false, true] {
let mut process = process();
let physical = process.db_startup.take().unwrap().into_process_capability();
let mut lease = admit(&process);
lease
.deadline
.timer
.as_mut()
.reset(tokio::time::Instant::now() + Duration::from_millis(5));
let mut scope = lease.into_serial_scope(std::future::pending());
{
let driver = scope.driver();
let result = driver
.supervise(async {
assert!(!panic_body, "injected business unwind");
loop {
driver.checkpoint().await?;
}
#[allow(unreachable_code)]
Ok::<(), ProfuseGwScopeStop>(())
})
.await;
if panic_body {
assert!(matches!(result, Err(ProfuseGwScopeFailure::Panicked)));
} else {
assert!(matches!(
result,
Err(ProfuseGwScopeFailure::Stopped(ProfuseGwScopeStop::TimedOut))
| Ok(Err(ProfuseGwScopeStop::TimedOut))
));
}
}
let (completion, execution) = scope.into_physical_finalization();
let physical = physical.connection_discarded(execution, ()).ok().unwrap();
let mut suspended = completion.complete(physical).ok().unwrap();
assert!(matches!(
suspended.resume().await,
Err(ProfuseGwScopeResumeError::Stopped(_))
));
let (_, terminal) = suspended.into_response_parts().ok().unwrap();
finish_profusegw_after_database(terminal);
process.finish().unwrap();
}
}
#[tokio::test]
async fn serial_scope_physical_expiry_and_restore_keep_stop() {
let mut process = process();
let physical = process.db_startup.take().unwrap().into_process_capability();
let mut lease = admit(&process);
let original_deadline = lease.deadline.timer.deadline();
lease.scope_stop = Some(ProfuseGwScopeStop::Cancelled);
let lease = lease.restore_dispatch().into_database_request();
assert_eq!(lease.scope_stop, Some(ProfuseGwScopeStop::Cancelled));
assert_eq!(lease.deadline.timer.deadline(), original_deadline);
let scope = lease.into_serial_scope(std::future::poll_fn(|_| -> Poll<()> {
panic!("sticky stop must not poll cancellation again")
}));
let (mut completion, execution) = scope.into_physical_finalization();
poll_fn(|cx| {
assert!(completion.poll_physical_stop(cx).is_ready());
Poll::Ready(())
})
.await;
let receipt = physical.connection_discarded(execution, ()).ok().unwrap();
let mut suspended = completion.complete(receipt).ok().unwrap();
assert!(matches!(
suspended.resume().await,
Err(ProfuseGwScopeResumeError::Stopped(
ProfuseGwScopeStop::Cancelled
))
));
let (_, terminal) = suspended.into_response_parts().ok().unwrap();
finish_profusegw_after_database(terminal);
process.finish().unwrap();
let mut process = self::process();
let physical = process.db_startup.take().unwrap().into_process_capability();
let scope = admit(&process).into_serial_scope(std::future::pending());
let (mut completion, execution) = scope.into_physical_finalization();
completion
.completion
.deadline
.timer
.as_mut()
.reset(tokio::time::Instant::now());
poll_fn(|cx| {
assert!(completion.poll_physical_stop(cx).is_ready());
Poll::Ready(())
})
.await;
let receipt = physical
.connection_returned(execution, "committed fact remains")
.ok()
.unwrap();
let mut suspended = completion.complete(receipt).ok().unwrap();
assert!(matches!(
suspended.resume().await,
Err(ProfuseGwScopeResumeError::Stopped(
ProfuseGwScopeStop::TimedOut
))
));
let (value, terminal) = suspended.into_response_parts().ok().unwrap();
assert_eq!(value, "committed fact remains");
finish_profusegw_after_database(terminal);
process.finish().unwrap();
}
}