mod observation;
pub use observation::{ReservedActiveStage, ReservedObservationStage};
pub(crate) mod scope;
pub use scope::{ReservedScopeCompletion, ReservedScopeObservation, ReservedScopeOperation, ReservedScopeSupervisionFacts, ReservedSupervisionFailure, ReservedSupervisedStage, supervised_scope_layouts};
use saddle_admission::{AdmissionError, StorageDemand, StoragePermit};
use saddle_core::request_context::{
ContextConflict, ContextFact, ContextLabel, RegisteredContextOperation, RequestExecutionView,
RequestIdentityGroup, RequestLocalFacts, RequestRootPublisher, RequestViewPhase,
};
use saddle_observability::EmergencyDiagnosticHandle;
use saddle_observability::root_diagnostic::{
PublicRequestFailure, RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent,
RootRequestFailure,
};
use std::{
alloc::Layout,
sync::{Arc, Mutex},
};
pub(crate) struct Retained<T>(Option<Arc<T>>);
impl<T> Retained<T> {
pub(crate) fn new(value: T) -> Self {
Self(Some(Arc::new(value)))
}
pub(crate) fn get_mut(&mut self) -> Option<&mut T> {
Arc::get_mut(self.0.as_mut()?)
}
}
impl<T> Clone for Retained<T> {
fn clone(&self) -> Self {
Self(Some(Arc::clone(
self.0.as_ref().expect("live shared storage"),
)))
}
}
impl<T> std::ops::Deref for Retained<T> {
type Target = Arc<T>;
fn deref(&self) -> &Arc<T> {
self.0.as_ref().expect("live shared storage")
}
}
impl<T> Drop for Retained<T> {
fn drop(&mut self) {
if let Some(shared) = self.0.take() {
drop(Arc::into_inner(shared));
}
}
}
#[derive(Debug)]
pub enum ReservedContextError {
Storage(AdmissionError),
Context(ContextConflict),
}
#[repr(C)]
struct RootStorage {
publisher: Mutex<RequestRootPublisher>,
permit: StoragePermit,
}
pub(crate) struct ReservedRequestLifetime(Retained<RootStorage>);
#[repr(C)]
struct ViewStorage {
view: RequestExecutionView,
child: Option<Retained<ChildStorage>>,
root: Retained<RootStorage>,
permit: StoragePermit,
}
#[repr(C)]
struct ChildStorage {
permit: StoragePermit,
}
pub struct ReservedRequestRoot(Retained<RootStorage>);
#[derive(Clone)]
pub struct ReservedRequestView(Retained<ViewStorage>);
pub struct ReservedTaskContext {
view: ReservedRequestView,
status: Retained<TaskStatusStorage>,
}
impl ReservedTaskContext {
pub fn view(&self) -> ReservedRequestView {
self.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view
.clone()
.unwrap_or_else(|| self.view.clone())
}
pub fn publish(&mut self, group: RequestIdentityGroup) -> Result<(), ReservedContextError> {
let [_, (_, layout), _] = saddle_core::request_context::request_context_layouts();
let permit = self
.view
.0
.permit
.try_reserve(view_demand(&[(layout, 1)]).map_err(ReservedContextError::Storage)?)
.map_err(ReservedContextError::Storage)?;
let mut publisher = self
.view
.0
.root
.publisher
.lock()
.unwrap_or_else(|p| p.into_inner());
publisher
.publish(group)
.map_err(|(error, _)| ReservedContextError::Context(error))?;
let view = self
.view
.0
.view
.refresh(&publisher.reference())
.map_err(ReservedContextError::Context)?;
drop(publisher);
self.view = ReservedRequestView(Retained::new(ViewStorage {
view,
child: self.view.0.child.clone(),
root: Retained::clone(&self.view.0.root),
permit,
}));
self.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view = Some(self.view.clone());
crate::diagnostics::update_root_view(&self.view);
Ok(())
}
pub fn observe_view(&mut self, view: ReservedRequestView) -> Result<(), ReservedRequestView> {
if !self.view.same_request(&view) {
return Err(view);
}
self.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view = Some(view.clone());
crate::diagnostics::update_root_view(&view);
self.view = view;
Ok(())
}
}
pub struct ReservedRequestFailure {
failure: RootRequestFailure,
source: ReservedRequestView,
}
fn shared_prefix<T>() -> Result<Layout, AdmissionError> {
Layout::new::<[std::sync::atomic::AtomicUsize; 2]>()
.extend(Layout::new::<T>())
.map(|(layout, _)| layout)
.map_err(|_| AdmissionError::SizeOverflow)
}
fn root_demand() -> Result<StorageDemand, AdmissionError> {
let [(_, root), _, _] = saddle_core::request_context::request_context_layouts();
StorageDemand::embedded(
shared_prefix::<Mutex<RequestRootPublisher>>()?,
Layout::new::<()>(),
&[(root, 1)],
)
}
fn dispatch_root_demand() -> Result<StorageDemand, AdmissionError> {
let [(_, root), _, _] = saddle_core::request_context::request_context_layouts();
StorageDemand::embedded(
shared_prefix::<Mutex<RequestRootPublisher>>()?,
Layout::new::<()>(),
&[(root, 1), (Layout::new::<tokio::time::Sleep>(), 1)],
)
}
fn view_demand(extra: &[(Layout, usize)]) -> Result<StorageDemand, AdmissionError> {
let prefix = shared_prefix::<(
RequestExecutionView,
Option<Retained<ChildStorage>>,
Retained<RootStorage>,
)>()?;
StorageDemand::embedded(prefix, Layout::new::<()>(), extra)
}
impl ReservedRequestRoot {
#[cfg(test)]
fn try_admitted(
original: saddle_admission::ProfuseGwLightweightPermit,
application: ContextLabel,
) -> Result<
(saddle_admission::ProfuseGwLightweightPermit, Self),
(
saddle_admission::ProfuseGwLightweightPermit,
ReservedContextError,
),
> {
let demand = match root_demand() {
Ok(demand) => demand,
Err(error) => return Err((original, ReservedContextError::Storage(error))),
};
let storage = match original.try_supervisor_storage(demand) {
Ok(storage) => storage,
Err(error) => return Err((original, ReservedContextError::Storage(error))),
};
match Self::create(storage, application) {
Ok(root) => Ok((original, root)),
Err(error) => Err((original, error)),
}
}
fn create(
permit: StoragePermit,
application: ContextLabel,
) -> Result<Self, ReservedContextError> {
let publisher = RequestRootPublisher::create(application, ContextFact::NotEstablished)
.map_err(ReservedContextError::Context)?;
Ok(Self(Retained::new(RootStorage {
publisher: Mutex::new(publisher),
permit,
})))
}
pub fn view(
&self,
phase: RequestViewPhase,
) -> Result<ReservedRequestView, ReservedContextError> {
let [_, (_, view), _] = saddle_core::request_context::request_context_layouts();
let permit = self
.0
.permit
.try_reserve(view_demand(&[(view, 1)]).map_err(ReservedContextError::Storage)?)
.map_err(ReservedContextError::Storage)?;
let view = self
.0
.publisher
.lock()
.unwrap_or_else(|p| p.into_inner())
.reference()
.view(RequestLocalFacts::new(phase));
Ok(ReservedRequestView(Retained::new(ViewStorage {
view,
child: None,
root: Retained::clone(&self.0),
permit,
})))
}
}
impl ReservedRequestView {
pub fn record_admission(
&self,
admission: crate::profusegw::ProfuseGwAdmissionEvent,
observer: &saddle_observability::Observer,
output: Option<&EmergencyDiagnosticHandle>,
construction: saddle_observability::AdmissionConstruction,
) -> saddle_observability::AdmissionCapacitySubmission {
admission.record_capacity_context(observer, output,
saddle_observability::AdmissionEventContext::Rooted(&self.0.view), construction)
}
#[cfg(test)]
pub(crate) fn test_available_storage(&self) -> usize {
let (mut low,mut high)=(0,65_536);
while low<high {
let mid=(low+high+1)/2;
let demand=StorageDemand::embedded(Layout::new::<()>(),Layout::new::<()>(),
&[(Layout::from_size_align(mid,1).unwrap(),1)]).unwrap();
match self.0.permit.try_reserve(demand) {
Ok(permit)=>{drop(permit);low=mid;}
Err(_)=>high=mid-1,
}
}
low
}
pub(crate) fn reserve_scope_body<F>(&self) -> Result<StoragePermit, AdmissionError> {
self.0.permit.try_reserve(StorageDemand::embedded(
Layout::new::<std::pin::Pin<Box<F>>>(),
Layout::new::<()>(),
&[(Layout::new::<F>(), 1)],
)?)
}
pub fn refresh(&self, root: &ReservedRequestRoot) -> Result<Self, ReservedContextError> {
if !Arc::ptr_eq(&self.0.root, &root.0) {
return Err(ReservedContextError::Context(ContextConflict::ForeignRoot));
}
let [_, (_, layout), _] = saddle_core::request_context::request_context_layouts();
let permit = self
.0
.permit
.try_reserve(view_demand(&[(layout, 1)]).map_err(ReservedContextError::Storage)?)
.map_err(ReservedContextError::Storage)?;
let view = self
.0
.view
.refresh(
&root
.0
.publisher
.lock()
.unwrap_or_else(|p| p.into_inner())
.reference(),
)
.map_err(ReservedContextError::Context)?;
Ok(Self(Retained::new(ViewStorage {
view,
child: self.0.child.clone(),
root: Retained::clone(&self.0.root),
permit,
})))
}
pub(crate) fn validate_execution(
&self,
execution: &saddle_admission::ProfuseGwLightweightExecutionOwner,
) -> Result<(), AdmissionError> {
execution.validate_supervisor_storage(&self.0.root.permit)
}
pub fn ordinary(
&self,
observer: &saddle_observability::Observer,
event: RootRequestEvent,
facts: RootOutcomeFacts,
) -> saddle_observability::DiagnosticSubmission {
RootDiagnosticScope::new(&self.0.view, None).ordinary(observer, event, facts)
}
fn in_actual_task(
&self,
id: tokio::task::Id,
permit: StoragePermit,
) -> Result<Self, ReservedContextError> {
struct Buffer {
bytes: [u8; 32],
len: usize,
}
impl std::fmt::Write for Buffer {
fn write_str(&mut self, s: &str) -> std::fmt::Result {
let end = self.len.checked_add(s.len()).ok_or(std::fmt::Error)?;
let target = self.bytes.get_mut(self.len..end).ok_or(std::fmt::Error)?;
target.copy_from_slice(s.as_bytes());
self.len = end;
Ok(())
}
}
use std::fmt::Write;
let mut buffer = Buffer {
bytes: [0; 32],
len: 0,
};
write!(&mut buffer, "{id}")
.map_err(|_| ReservedContextError::Context(ContextConflict::InvalidIdentity))?;
let number = std::str::from_utf8(&buffer.bytes[..buffer.len])
.ok()
.and_then(|s| s.parse::<u64>().ok())
.ok_or(ReservedContextError::Context(
ContextConflict::InvalidIdentity,
))?;
Ok(Self(Retained::new(ViewStorage {
view: self.0.view.in_task(number),
child: self.0.child.clone(),
root: Retained::clone(&self.0.root),
permit,
})))
}
pub(crate) fn in_scope(
&self,
request: &saddle_core::DbPhysicalRequestHalf,
execution: &saddle_core::DbPhysicalExecutionHalf,
) -> Result<Self, ReservedContextError> {
let checked = request
.project_diagnostic_context(execution, self.0.view.clone())
.map_err(|_| ReservedContextError::Context(ContextConflict::ForeignRoot))?;
let [_, (_, layout), _] = saddle_core::request_context::request_context_layouts();
let permit = self
.0
.permit
.try_reserve(view_demand(&[(layout, 1)]).map_err(ReservedContextError::Storage)?)
.map_err(ReservedContextError::Storage)?;
let view = self
.0
.view
.in_db_scope(&checked)
.map_err(ReservedContextError::Context)?;
Ok(Self(Retained::new(ViewStorage {
view,
child: self.0.child.clone(),
root: Retained::clone(&self.0.root),
permit,
})))
}
pub fn child(
&self,
call: &saddle_core::CallContext,
request: &str,
route: &str,
attempt: u32,
) -> Result<Self, ReservedContextError> {
let [_, (_, view_layout), (_, child_layout)] =
saddle_core::request_context::request_context_layouts();
let permit = self
.0
.permit
.try_reserve(view_demand(&[(view_layout, 1)]).map_err(ReservedContextError::Storage)?)
.map_err(ReservedContextError::Storage)?;
let child_permit = self
.0
.permit
.try_reserve(
StorageDemand::embedded(
Layout::new::<[std::sync::atomic::AtomicUsize; 2]>(),
Layout::new::<()>(),
&[(child_layout, 1)],
)
.map_err(ReservedContextError::Storage)?,
)
.map_err(ReservedContextError::Storage)?;
let view = self
.0
.view
.child(call, request, route, attempt)
.map_err(ReservedContextError::Context)?;
Ok(Self(Retained::new(ViewStorage {
view,
child: Some(Retained::new(ChildStorage {
permit: child_permit,
})),
root: Retained::clone(&self.0.root),
permit,
})))
}
pub(crate) fn panic_source(
&self,
info: &std::panic::PanicHookInfo<'_>,
diagnostic: saddle_core::Diagnostic,
output: Option<&EmergencyDiagnosticHandle>,
) -> ReservedRequestFailure {
let description = info
.payload()
.downcast_ref::<&str>()
.copied()
.or_else(|| info.payload().downcast_ref::<String>().map(String::as_str))
.unwrap_or("panic payload does not expose a string description");
self.existing_description(description, diagnostic, output)
}
pub(crate) fn existing_description(
&self,
description: &str,
diagnostic: saddle_core::Diagnostic,
output: Option<&EmergencyDiagnosticHandle>,
) -> ReservedRequestFailure {
ReservedRequestFailure {
failure: RootDiagnosticScope::new(&self.0.view, output).source_existing_description(
&description,
diagnostic,
saddle_core::DiagnosticCode::new("runtime.request_task_failure").unwrap(),
RootRequestEvent::Finalization,
RootOutcomeFacts::default(),
),
source: self.clone(),
}
}
fn derive(
&self,
build: impl FnOnce(&RequestExecutionView) -> RequestExecutionView,
) -> Result<Self, ReservedContextError> {
let [_, (_, view), _] = saddle_core::request_context::request_context_layouts();
let permit = self
.0
.permit
.try_reserve(view_demand(&[(view, 1)]).map_err(ReservedContextError::Storage)?)
.map_err(ReservedContextError::Storage)?;
Ok(Self(Retained::new(ViewStorage {
view: build(&self.0.view),
child: self.0.child.clone(),
root: Retained::clone(&self.0.root),
permit,
})))
}
pub fn with_phase(&self, phase: RequestViewPhase) -> Result<Self, ReservedContextError> {
self.derive(|view| view.with_phase(phase))
}
pub fn with_db_operation(
&self,
operation: RegisteredContextOperation,
) -> Result<Self, ReservedContextError> {
self.derive(|view| view.with_db_operation(operation))
}
pub fn same_request(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.0.root, &other.0.root)
}
#[track_caller]
pub fn source_error(
&self,
error: &(dyn std::error::Error + 'static),
output: Option<&EmergencyDiagnosticHandle>,
stage: saddle_core::DiagnosticStage,
event: RootRequestEvent,
facts: RootOutcomeFacts,
) -> ReservedRequestFailure {
ReservedRequestFailure {
failure: RootDiagnosticScope::new(&self.0.view, output)
.source_error(error, stage, event, facts),
source: self.clone(),
}
}
}
impl ReservedRequestFailure {
pub fn occurrence(&self) -> saddle_core::DiagnosticOccurrence {
self.failure.occurrence()
}
pub fn finish(
self,
current: &ReservedRequestView,
output: Option<&EmergencyDiagnosticHandle>,
facts: RootOutcomeFacts,
) -> Result<PublicRequestFailure, Self> {
if !self.source.same_request(current) {
return Err(self);
}
let Self { failure, source } = self;
match failure.finish(¤t.0.view, output, facts) {
Ok(public) => Ok(public),
Err(failure) => Err(Self { failure, source }),
}
}
}
pub type ReservedBorrowedFuture<'a, T> = std::pin::Pin<
Box<dyn std::future::Future<Output = Result<T, ReservedRequestFailure>> + Send + 'a>,
>;
pub fn dispatch_storage_bytes_for<X: Send + 'static, T: Send + 'static, F>(
_: &F,
body: Layout,
indirect: &[(Layout, usize)],
) -> Result<usize, AdmissionError>
where
F: for<'a> FnOnce(
&'a mut crate::profusegw::ReservedDispatchOwner<X>,
ReservedTaskContext,
) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
let [_, (_, view), _] = saddle_core::request_context::request_context_layouts();
let parts = [
dispatch_root_demand()?.bytes(),
view_demand(&[(view, 1)])?.bytes(),
task_demand_with::<crate::profusegw::ReservedDispatchOwner<X>, T, F>(
body, indirect,
)?
.bytes(),
view_demand(&[(view, 1)])?.bytes(),
];
parts.into_iter().try_fold(0usize, |sum, part| {
sum.checked_add(part).ok_or(AdmissionError::SizeOverflow)
})
}
struct TaskStatus {
id: Option<tokio::task::Id>,
binding_failed: bool,
primary: Option<ReservedRequestFailure>,
cleanup: Option<ReservedRequestFailure>,
scopes: Option<Retained<scope::ScopeStorage>>,
task_view: Option<ReservedRequestView>,
task_view_permit: Option<StoragePermit>,
}
#[repr(C)]
struct TaskStatusStorage {
status: Mutex<TaskStatus>,
permit: StoragePermit,
}
impl std::ops::Deref for TaskStatusStorage {
type Target = Mutex<TaskStatus>;
fn deref(&self) -> &Self::Target {
&self.status
}
}
#[repr(C)]
struct TaskContent<O> {
owner: Arc<tokio::sync::Mutex<Option<O>>>,
view: ReservedRequestView,
output: Option<EmergencyDiagnosticHandle>,
status: Retained<TaskStatusStorage>,
}
#[repr(C)]
struct TaskStorage<O> {
content: TaskContent<O>,
}
pub struct ReservedTaskTicket<O>(Retained<TaskStorage<O>>);
pub struct ReservedTaskOutput<O, T> {
result: Option<Result<T, ReservedRequestFailure>>,
shared: Retained<TaskStorage<O>>,
}
pub struct ReservedTaskFuture<O, T> {
body: Option<
std::pin::Pin<
Box<dyn std::future::Future<Output = Option<Result<T, ReservedRequestFailure>>> + Send>,
>,
>,
shared: Retained<TaskStorage<O>>,
completed: bool,
}
async fn borrowed_body<O: Send + 'static, T>(
shared: Retained<TaskStorage<O>>,
make: impl for<'a> FnOnce(&'a mut O, ReservedTaskContext) -> ReservedBorrowedFuture<'a, T> + Send,
) -> Option<Result<T, ReservedRequestFailure>> {
let mut owner = shared
.content
.owner
.clone()
.try_lock_owned()
.expect("one internal borrower");
let view = shared
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view
.clone()
.expect("poll establishes task view");
let mut body = make(
owner.as_mut().expect("owner before join"),
ReservedTaskContext {
view: view.clone(),
status: Retained::clone(&shared.content.status),
},
);
let result = std::future::poll_fn(|cx| {
let view = shared
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view
.clone()
.expect("established view");
match crate::diagnostics::catching_root(
view,
shared.content.output.clone(),
None,
false,
|| body.as_mut().poll(cx),
) {
Ok(std::task::Poll::Pending) => std::task::Poll::Pending,
Ok(std::task::Poll::Ready(result)) => std::task::Poll::Ready(Some(result)),
Err(failure) => {
shared
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.primary = Some(failure);
std::task::Poll::Ready(None)
}
}
})
.await;
let primary = result
.as_ref()
.and_then(|r| r.as_ref().err())
.map(ReservedRequestFailure::occurrence)
.or_else(|| {
shared
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.primary
.as_ref()
.map(ReservedRequestFailure::occurrence)
});
let cleanup_view = shared
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view
.clone()
.unwrap_or(view);
if let Err(cleanup) = crate::diagnostics::catching_root(
cleanup_view,
shared.content.output.clone(),
primary,
true,
|| drop(body),
) {
shared
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.cleanup = Some(cleanup);
}
result
}
fn return_layout<I, R>(_: impl FnOnce(I) -> R) -> Layout {
Layout::new::<R>()
}
#[cfg(test)]
fn task_demand<O: Send + 'static, T: Send + 'static, F>(
body: Layout,
) -> Result<StorageDemand, AdmissionError>
where
F: for<'a> FnOnce(&'a mut O, ReservedTaskContext) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
task_demand_with::<O, T, F>(body, &[])
}
fn task_demand_with<O: Send + 'static, T: Send + 'static, F>(
body: Layout,
indirect: &[(Layout, usize)],
) -> Result<StorageDemand, AdmissionError>
where
F: for<'a> FnOnce(&'a mut O, ReservedTaskContext) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
if !cfg!(all(
target_arch = "x86_64",
target_os = "linux",
target_pointer_width = "64"
)) {
return Err(AdmissionError::InvalidConfiguration);
}
let future =
return_layout(|(shared, make): (Retained<TaskStorage<O>>, F)| borrowed_body(shared, make));
let owner_cell = super::storage::shared(Layout::new::<tokio::sync::Mutex<Option<O>>>())
.map_err(|_| AdmissionError::SizeOverflow)?
.allocation;
let abi = super::storage::TaskBackendLayouts {
header: Layout::from_size_align(32, 8).unwrap(),
scheduler: Layout::new::<usize>(),
task_id: Layout::new::<u64>(),
trailer: Layout::from_size_align(48, 8).unwrap(),
cell_alignment: 128,
joinset_entry: Layout::from_size_align(56, 8).unwrap(),
joinset_shared: Layout::from_size_align(72, 8).unwrap(),
};
let cell = super::storage::task_cell(
abi,
Layout::new::<ReservedTaskFuture<O, T>>(),
Layout::new::<Result<ReservedTaskOutput<O, T>, tokio::task::JoinError>>(),
)
.map_err(|_| AdmissionError::SizeOverflow)?;
let indirect_bytes = indirect.iter().try_fold(0usize, |sum, (layout, count)| {
layout
.pad_to_align()
.size()
.checked_mul(*count)
.and_then(|n| sum.checked_add(n))
.ok_or(AdmissionError::SizeOverflow)
})?;
let indirect =
Layout::from_size_align(indirect_bytes, 1).map_err(|_| AdmissionError::SizeOverflow)?;
let task_shared = super::storage::shared(Layout::new::<TaskStorage<O>>())
.map_err(|_| AdmissionError::SizeOverflow)?
.allocation;
StorageDemand::embedded(
shared_prefix::<Mutex<TaskStatus>>()?,
Layout::new::<()>(),
&[
(task_shared, 1),
(owner_cell, 1),
(future, 1),
(body, 1),
(cell, 1),
(abi.joinset_entry, 1),
(indirect, 1),
],
)
}
pub(crate) struct ReservedTaskAllocation {
root: StoragePermit,
initial: StoragePermit,
task: StoragePermit,
task_view: StoragePermit,
}
#[cfg(test)]
std::thread_local! {
pub(crate) static FAIL_AR_CONSTRUCTION: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
}
impl ReservedTaskAllocation {
pub(crate) fn for_dispatch<X: Send + 'static, T: Send + 'static, F>(
permit: &saddle_admission::ProfuseGwLightweightPermit,
body: Layout,
indirect: &[(Layout, usize)],
) -> Result<Self, AdmissionError>
where
F: for<'a> FnOnce(
&'a mut crate::profusegw::ReservedDispatchOwner<X>,
ReservedTaskContext,
) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
let root = permit.try_supervisor_storage(dispatch_root_demand()?)?;
let [_, (_, view), _] = saddle_core::request_context::request_context_layouts();
let initial = root.try_reserve(view_demand(&[(view, 1)])?)?;
let task = root.try_reserve(task_demand_with::<
crate::profusegw::ReservedDispatchOwner<X>,
T,
F,
>(body, indirect)?)?;
let task_view = root.try_reserve(view_demand(&[(view, 1)])?)?;
Ok(Self {
root,
initial,
task,
task_view,
})
}
pub(crate) fn construct<X: Send + 'static, T: Send + 'static, F>(
self,
mut dispatch: crate::profusegw::ProfuseGwManagedDispatch,
input: X,
admission: crate::profusegw::ProfuseGwAdmissionEvent,
application: ContextLabel,
output: Option<EmergencyDiagnosticHandle>,
make: F,
) -> Result<
(
ReservedRequestRoot,
ReservedTaskFuture<crate::profusegw::ReservedDispatchOwner<X>, T>,
ReservedTaskTicket<crate::profusegw::ReservedDispatchOwner<X>>,
),
(
crate::profusegw::ProfuseGwManagedDispatch,
X,
crate::profusegw::ProfuseGwAdmissionEvent,
F,
ReservedContextError,
),
>
where
F: for<'a> FnOnce(
&'a mut crate::profusegw::ReservedDispatchOwner<X>,
ReservedTaskContext,
) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
#[cfg(test)]
if FAIL_AR_CONSTRUCTION.with(|flag| flag.replace(false)) {
return Err((dispatch, input, admission, make,
ReservedContextError::Storage(AdmissionError::AccountClosed)));
}
let root = match ReservedRequestRoot::create(self.root, application) {
Ok(root) => root,
Err(error) => return Err((dispatch, input, admission, make, error)),
};
dispatch.retain_reserved_deadline(ReservedRequestLifetime(Retained::clone(&root.0)));
let raw = root
.0
.publisher
.lock()
.unwrap_or_else(|p| p.into_inner())
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Reading));
let view = ReservedRequestView(Retained::new(ViewStorage {
view: raw,
child: None,
root: Retained::clone(&root.0),
permit: self.initial,
}));
let (future, ticket) = build_task(
(Some(dispatch), input, Some(admission)),
view,
output,
self.task,
self.task_view,
make,
);
Ok((root, future, ticket))
}
}
fn build_task<O: Send + 'static, T: Send + 'static, F>(
owner: O,
view: ReservedRequestView,
output: Option<EmergencyDiagnosticHandle>,
permit: StoragePermit,
task_view_permit: StoragePermit,
make: F,
) -> (ReservedTaskFuture<O, T>, ReservedTaskTicket<O>)
where
F: for<'a> FnOnce(&'a mut O, ReservedTaskContext) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
let owner = Arc::new(tokio::sync::Mutex::new(Some(owner)));
let status = Retained::new(TaskStatusStorage {
status: Mutex::new(TaskStatus {
id: None,
binding_failed: false,
primary: None,
cleanup: None,
scopes: None,
task_view: None,
task_view_permit: Some(task_view_permit),
}),
permit,
});
let shared = Retained::new(TaskStorage {
content: TaskContent {
owner,
view,
output,
status,
},
});
let inner = borrowed_body(Retained::clone(&shared), make);
(
ReservedTaskFuture {
body: Some(Box::pin(inner)),
shared: Retained::clone(&shared),
completed: false,
},
ReservedTaskTicket(shared),
)
}
#[cfg(test)]
fn try_reserved_task<O: Send + 'static, T: Send + 'static, F>(
owner: O,
view: ReservedRequestView,
output: Option<EmergencyDiagnosticHandle>,
body: Layout,
make: F,
) -> Result<
(ReservedTaskFuture<O, T>, ReservedTaskTicket<O>),
(O, ReservedRequestView, AdmissionError),
>
where
F: for<'a> FnOnce(&'a mut O, ReservedTaskContext) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
let demand = match task_demand::<O, T, F>(body) {
Ok(d) => d,
Err(e) => return Err((owner, view, e)),
};
let permit = match view.0.permit.try_reserve(demand) {
Ok(p) => p,
Err(e) => return Err((owner, view, e)),
};
let [_, (_, local), _] = saddle_core::request_context::request_context_layouts();
let task_view_demand = match view_demand(&[(local, 1)]) {
Ok(d) => d,
Err(e) => return Err((owner, view, e)),
};
let task_view_permit = match view.0.permit.try_reserve(task_view_demand) {
Ok(p) => p,
Err(e) => return Err((owner, view, e)),
};
Ok(build_task(
owner,
view,
output,
permit,
task_view_permit,
make,
))
}
#[cfg(test)]
fn try_prepare_execution_task<T: Send + 'static, F>(
original: saddle_admission::ProfuseGwLightweightPermit,
application: ContextLabel,
output: Option<EmergencyDiagnosticHandle>,
body: Layout,
make: F,
) -> Result<
(
ReservedRequestRoot,
ReservedTaskFuture<saddle_admission::ProfuseGwLightweightExecutionOwner, T>,
ReservedTaskTicket<saddle_admission::ProfuseGwLightweightExecutionOwner>,
),
(
saddle_admission::ProfuseGwLightweightPermit,
F,
ReservedContextError,
),
>
where
F: for<'a> FnOnce(
&'a mut saddle_admission::ProfuseGwLightweightExecutionOwner,
ReservedTaskContext,
) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
let reserve = || -> Result<_, AdmissionError> {
let root = original.try_supervisor_storage(root_demand()?)?;
let [_, (_, view), _] = saddle_core::request_context::request_context_layouts();
let initial = root.try_reserve(view_demand(&[(view, 1)])?)?;
let task = root.try_reserve(task_demand::<
saddle_admission::ProfuseGwLightweightExecutionOwner,
T,
F,
>(body)?)?;
let task_view = root.try_reserve(view_demand(&[(view, 1)])?)?;
Ok((root, initial, task, task_view))
};
let (root, initial, task, task_view) = match reserve() {
Ok(p) => p,
Err(e) => return Err((original, make, ReservedContextError::Storage(e))),
};
let root = match ReservedRequestRoot::create(root, application) {
Ok(r) => r,
Err(e) => return Err((original, make, e)),
};
let raw = root
.0
.publisher
.lock()
.unwrap_or_else(|p| p.into_inner())
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Reading));
let view = ReservedRequestView(Retained::new(ViewStorage {
view: raw,
child: None,
root: Retained::clone(&root.0),
permit: initial,
}));
let (future, ticket) = build_task(
original.into_execution(),
view,
output,
task,
task_view,
make,
);
Ok((root, future, ticket))
}
pub(crate) fn try_prepare_rejection_task<O: Send + 'static, T: Send + 'static, F>(
profile: &saddle_admission::VerifiedProfuseGwLightweightProfile<'_>,
application: ContextLabel,
output: Option<EmergencyDiagnosticHandle>,
body: Layout,
indirect: &[(Layout, usize)],
owner: O,
make: F,
) -> Result<
(
ReservedRequestRoot,
ReservedTaskFuture<O, T>,
ReservedTaskTicket<O>,
),
(O, F, ReservedContextError),
>
where
F: for<'a> FnOnce(&'a mut O, ReservedTaskContext) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
let reserve = || -> Result<_, AdmissionError> {
let root = profile.try_rejection_storage(root_demand()?)?;
let [_, (_, view), _] = saddle_core::request_context::request_context_layouts();
let initial = root.try_reserve(view_demand(&[(view, 1)])?)?;
let task = root.try_reserve(task_demand_with::<O, T, F>(body, indirect)?)?;
let task_view = root.try_reserve(view_demand(&[(view, 1)])?)?;
Ok((root, initial, task, task_view))
};
let (root, initial, task, task_view) = match reserve() {
Ok(p) => p,
Err(e) => return Err((owner, make, ReservedContextError::Storage(e))),
};
let root = match ReservedRequestRoot::create(root, application) {
Ok(r) => r,
Err(e) => return Err((owner, make, e)),
};
let raw = root
.0
.publisher
.lock()
.unwrap_or_else(|p| p.into_inner())
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Reading));
let view = ReservedRequestView(Retained::new(ViewStorage {
view: raw,
child: None,
root: Retained::clone(&root.0),
permit: initial,
}));
let (future, ticket) = build_task(owner, view, output, task, task_view, make);
Ok((root, future, ticket))
}
impl<O, T> Unpin for ReservedTaskFuture<O, T> {}
impl<O, T> std::future::Future for ReservedTaskFuture<O, T> {
type Output = ReservedTaskOutput<O, T>;
fn poll(
self: std::pin::Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Self::Output> {
let this = self.get_mut();
let content = &this.shared.content;
let Some(id) = tokio::task::try_id() else {
return this.reject_binding();
};
{
let mut status = content.status.lock().unwrap_or_else(|p| p.into_inner());
if status.binding_failed || status.id.is_some_and(|old| old != id) {
drop(status);
return this.reject_binding();
}
status.id = Some(id);
if status.task_view.is_none() {
let Some(prepaid) = status.task_view_permit.take() else {
drop(status);
return this.reject_binding();
};
match content.view.in_actual_task(id, prepaid) {
Ok(view) => status.task_view = Some(view),
Err(_) => {
status.primary = Some(content.view.existing_description(
"task context storage unavailable",
crate::diagnostics::failure(
saddle_core::DiagnosticStage::BackgroundTask,
saddle_core::DiagnosticCategory::ExpectedRejection,
"runtime.task_context_storage_unavailable",
),
content.output.as_ref(),
));
drop(status);
this.drop_body(None);
this.completed = true;
return std::task::Poll::Ready(ReservedTaskOutput {
result: None,
shared: Retained::clone(&this.shared),
});
}
}
}
}
let view = content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view
.clone()
.expect("task view");
let outcome =
crate::diagnostics::catching_root(view, content.output.clone(), None, false, || {
this.body
.as_mut()
.expect("no poll after Ready")
.as_mut()
.poll(cx)
});
let result = match outcome {
Ok(std::task::Poll::Pending) => return std::task::Poll::Pending,
Ok(std::task::Poll::Ready(result)) => result,
Err(failure) => {
content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.primary = Some(failure);
None
}
};
let primary = result
.as_ref()
.and_then(|r| r.as_ref().err())
.map(ReservedRequestFailure::occurrence);
this.drop_body(primary);
this.completed = true;
std::task::Poll::Ready(ReservedTaskOutput {
result,
shared: Retained::clone(&this.shared),
})
}
}
impl<O, T> ReservedTaskFuture<O, T> {
fn reject_binding(&mut self) -> std::task::Poll<ReservedTaskOutput<O, T>> {
let content = &self.shared.content;
let failure = content.view.existing_description(
"request task identity mismatch or unavailable",
crate::diagnostics::failure(
saddle_core::DiagnosticStage::BackgroundTask,
saddle_core::DiagnosticCategory::ExpectedRejection,
"runtime.task_identity_unavailable",
),
content.output.as_ref(),
);
{
let mut s = content.status.lock().unwrap_or_else(|p| p.into_inner());
s.binding_failed = true;
s.primary.get_or_insert(failure);
}
self.drop_body(None);
self.completed = true;
std::task::Poll::Ready(ReservedTaskOutput {
result: None,
shared: Retained::clone(&self.shared),
})
}
fn drop_body(&mut self, primary: Option<saddle_core::DiagnosticOccurrence>) {
let Some(body) = self.body.take() else { return };
let content = &self.shared.content;
let primary = primary.or_else(|| {
content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.primary
.as_ref()
.map(ReservedRequestFailure::occurrence)
});
let view = content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view
.clone()
.unwrap_or_else(|| content.view.clone());
if let Err(cleanup) =
crate::diagnostics::catching_root(view, content.output.clone(), primary, true, || {
drop(body)
})
{
content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.cleanup = Some(cleanup);
}
}
}
impl<O, T> Drop for ReservedTaskFuture<O, T> {
fn drop(&mut self) {
if !self.completed {
let content = &self.shared.content;
let view = content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view
.clone()
.unwrap_or_else(|| content.view.clone());
let failure = view.existing_description(
"request task cancelled",
crate::diagnostics::failure(
saddle_core::DiagnosticStage::BackgroundTask,
saddle_core::DiagnosticCategory::ExpectedRejection,
"runtime.scope_cancelled",
),
content.output.as_ref(),
);
content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.primary
.get_or_insert(failure);
}
self.drop_body(None);
}
}
pub struct ReservedTaskJoined<O, T> {
result: Option<Result<T, ReservedRequestFailure>>,
primary: Option<ReservedRequestFailure>,
cleanup: Option<ReservedRequestFailure>,
owner: O,
storage: Retained<TaskStorage<O>>,
}
pub struct ReservedTaskRecovered<O, T> {
pub result: Option<Result<T, PublicRequestFailure>>,
pub primary: Option<PublicRequestFailure>,
pub cleanup: Option<PublicRequestFailure>,
pub owner: O,
}
impl<O, T> ReservedTaskJoined<O, T> {
pub fn record_admission(
&self,
admission: crate::profusegw::ProfuseGwAdmissionEvent,
observer: &saddle_observability::Observer,
construction: saddle_observability::AdmissionConstruction,
) -> saddle_observability::AdmissionCapacitySubmission {
let current = self.storage.content.status.lock()
.unwrap_or_else(|p| p.into_inner()).task_view.clone()
.unwrap_or_else(|| self.storage.content.view.clone());
current.record_admission(admission, observer,
self.storage.content.output.as_ref(), construction)
}
pub fn owner_mut(&mut self) -> &mut O {
&mut self.owner
}
pub fn recover(self, mut facts: RootOutcomeFacts) -> Result<ReservedTaskRecovered<O, T>, Self> {
let current = self
.storage
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.task_view
.clone()
.unwrap_or_else(|| self.storage.content.view.clone());
let failures = self
.result
.as_ref()
.and_then(|result| result.as_ref().err())
.into_iter()
.chain(self.primary.as_ref())
.chain(self.cleanup.as_ref());
if failures
.into_iter()
.any(|failure| !failure.source.same_request(¤t))
{
return Err(self);
}
let Self {
result,
primary,
cleanup,
owner,
storage,
} = self;
let scopes = storage.content.status.lock().unwrap_or_else(|p| p.into_inner()).scopes.take();
if cleanup.is_some() || scope::has_cleanup(&scopes) {
facts.axes.cleanup = saddle_core::CleanupOutcome::Failed;
}
let project = |failure: ReservedRequestFailure| {
let ReservedRequestFailure { failure, source } = failure;
let _source_guard = source;
match failure.finish(¤t.0.view, storage.content.output.as_ref(), facts) {
Ok(public) => public,
Err(_) => unreachable!("same source projection"),
}
};
let result = result.map(|result| result.map_err(&project));
let mut primary = primary.map(&project);
let mut cleanup = cleanup.map(&project);
let (scope_primary, scope_cleanup) = scope::project_scopes(scopes, &project);
primary = primary.or(scope_primary);
cleanup = cleanup.or(scope_cleanup);
drop(current);
drop(storage);
Ok(ReservedTaskRecovered {
result,
primary,
cleanup,
owner,
})
}
pub fn take_foreign_failure(&mut self) -> Option<ReservedRequestFailure> {
let foreign = self
.result
.as_ref()
.and_then(|r| r.as_ref().err())
.is_some_and(|f| !f.source.same_request(&self.storage.content.view));
if !foreign {
return None;
}
let Some(Err(failure)) = self.result.take() else {
return None;
};
self.primary = Some(self.storage.content.view.existing_description(
"foreign request failure receipt",
crate::diagnostics::failure(
saddle_core::DiagnosticStage::BackgroundTask,
saddle_core::DiagnosticCategory::ExpectedRejection,
"runtime.foreign_failure_receipt",
),
self.storage.content.output.as_ref(),
));
Some(failure)
}
}
impl<O> ReservedTaskTicket<O> {
pub fn withdraw<T>(
self,
future: ReservedTaskFuture<O, T>,
) -> Result<ReservedTaskJoined<O, T>, (Self, ReservedTaskFuture<O, T>)> {
if !self.matches_future(&future)
|| self
.0
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.id
.is_some()
{
return Err((self, future));
}
drop(future); let owner = self
.0
.content
.owner
.try_lock()
.ok()
.and_then(|mut owner| owner.take());
let Some(owner) = owner else {
unreachable!("unsubmitted future has released the only borrower")
};
let (primary, cleanup) = {
let mut s = self
.0
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner());
(s.primary.take(), s.cleanup.take())
};
Ok(ReservedTaskJoined {
result: None,
primary,
cleanup,
owner,
storage: self.0,
})
}
pub(super) fn matches_future<T>(&self, future: &ReservedTaskFuture<O, T>) -> bool {
Arc::ptr_eq(&self.0, &future.shared)
}
pub fn bind(&mut self, id: tokio::task::Id) -> Result<(), ()> {
let mut status = self
.0
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner());
if status.id.is_some_and(|old| old != id) {
return Err(());
}
status.id = Some(id);
if status.task_view.is_none() {
if let Some(permit) = status.task_view_permit.take() {
match self.0.content.view.in_actual_task(id, permit) {
Ok(view) => status.task_view = Some(view),
Err(_) => {
status.binding_failed = true;
return Err(());
}
}
}
}
Ok(())
}
pub fn complete<T>(
self,
result: Result<ReservedTaskOutput<O, T>, tokio::task::JoinError>,
) -> Result<
ReservedTaskJoined<O, T>,
(
Self,
Result<ReservedTaskOutput<O, T>, tokio::task::JoinError>,
),
> {
let matches = match &result {
Ok(output) => Arc::ptr_eq(&self.0, &output.shared),
Err(error) => {
self.0
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner())
.id
== Some(error.id())
}
};
if !matches {
return Err((self, result));
}
let owner = self
.0
.content
.owner
.try_lock()
.ok()
.and_then(|mut owner| owner.take());
let Some(owner) = owner else {
return Err((self, result));
};
let result = result.ok().and_then(|output| output.result);
let (primary, cleanup) = {
let mut s = self
.0
.content
.status
.lock()
.unwrap_or_else(|p| p.into_inner());
(s.primary.take(), s.cleanup.take())
};
Ok(ReservedTaskJoined {
result,
primary,
cleanup,
owner,
storage: self.0,
})
}
}
#[cfg(test)]
fn try_reserved_execution_task<T: Send + 'static, F>(
owner: saddle_admission::ProfuseGwLightweightExecutionOwner,
view: ReservedRequestView,
output: Option<EmergencyDiagnosticHandle>,
body: Layout,
make: F,
) -> Result<
(
ReservedTaskFuture<saddle_admission::ProfuseGwLightweightExecutionOwner, T>,
ReservedTaskTicket<saddle_admission::ProfuseGwLightweightExecutionOwner>,
),
(
saddle_admission::ProfuseGwLightweightExecutionOwner,
ReservedRequestView,
AdmissionError,
),
>
where
F: for<'a> FnOnce(
&'a mut saddle_admission::ProfuseGwLightweightExecutionOwner,
ReservedTaskContext,
) -> ReservedBorrowedFuture<'a, T>
+ Send
+ 'static,
{
if let Err(error) = view.validate_execution(&owner) {
return Err((owner, view, error));
}
try_reserved_task(owner, view, output, body, make)
}
#[cfg(test)]
pub(crate) mod tests {
use super::*;
pub(crate) fn process() -> saddle_admission::ProfuseGwLightweightProcessOwner {
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(),
std::time::Duration::from_millis(5000),
))
.ok()
.unwrap();
let (whole, receipt) = saddle_core::pair_bootstrap_rendezvous(app, listener)
.ok()
.unwrap();
let budget =
saddle_admission::bind_deployment_resource_budget_bootstrap(pending, whole, receipt)
.ok()
.unwrap();
saddle_admission::prepare_profusegw_lightweight_profile(budget)
.ok()
.unwrap()
}
fn admit(
process: &saddle_admission::ProfuseGwLightweightProcessOwner,
) -> saddle_admission::ProfuseGwLightweightPermit {
match process.verified_profile().try_admit() {
saddle_admission::ProfuseGwLightweightAdmissionOutcome::Ready(p) => p,
_ => panic!("original capacity"),
}
}
fn rejected(process: &saddle_admission::ProfuseGwLightweightProcessOwner) {
assert!(matches!(
process.verified_profile().try_admit(),
saddle_admission::ProfuseGwLightweightAdmissionOutcome::CapacityRejected
));
}
#[test]
fn real_storage_tracks_last_guarded_view_and_public_error_is_detached() {
let process = process();
let original = admit(&process);
let (original, root) =
ReservedRequestRoot::try_admitted(original, ContextLabel::checked("app").unwrap())
.ok()
.unwrap();
let view = root.view(RequestViewPhase::Reading).unwrap();
let held = view.clone();
let weak = Arc::downgrade(&root.0);
let failure = view.source_error(
&std::io::Error::other("original error"),
None,
saddle_core::DiagnosticStage::RequestDb,
RootRequestEvent::Finalization,
RootOutcomeFacts::default(),
);
let public = failure
.finish(&view, None, RootOutcomeFacts::default())
.ok()
.unwrap();
let execution = original.into_execution();
let terminal = execution.cancel_observed();
assert_eq!(
terminal.post_release().cpu_used(),
1,
"storage still retained"
);
drop((root, view));
rejected(&process);
assert!(weak.upgrade().is_some());
drop(held);
assert!(weak.upgrade().is_none());
drop(admit(&process));
std::hint::black_box(public);
process.finish().unwrap();
}
#[test]
fn matching_join_retains_owner_and_storage_across_pre_poll_pending_and_ready() {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(async {
for stage in 0..3 {
let process = process();
let (permit, root) = ReservedRequestRoot::try_admitted(
admit(&process),
ContextLabel::checked("app").unwrap(),
)
.ok()
.unwrap();
let view = root.view(RequestViewPhase::Reading).unwrap();
let entered = Arc::new(std::sync::atomic::AtomicBool::new(false));
let observed = entered.clone();
let body =
Layout::new::<std::future::Ready<Result<u32, ReservedRequestFailure>>>();
let (future, mut ticket) = try_reserved_execution_task(
permit.into_execution(),
view,
None,
body,
move |_, _| {
observed.store(true, std::sync::atomic::Ordering::SeqCst);
if stage == 2 {
Box::pin(std::future::ready(Ok(41)))
} else {
Box::pin(std::future::pending())
}
},
)
.ok()
.expect("real storage reservation fits isolated consumer");
let weak = Arc::downgrade(&ticket.0);
let mut set = tokio::task::JoinSet::new();
let handle = set.spawn(future);
ticket.bind(handle.id()).unwrap();
if stage != 0 {
while !entered.load(std::sync::atomic::Ordering::SeqCst) {
tokio::task::yield_now().await;
}
tokio::task::yield_now().await;
}
if stage != 2 {
handle.abort();
}
assert_eq!(
entered.load(std::sync::atomic::Ordering::SeqCst),
stage != 0
);
rejected(&process);
assert!(
weak.upgrade().is_some(),
"body completion does not refund ticket"
);
let joined = ticket
.complete(set.join_next().await.unwrap())
.ok()
.unwrap();
drop(root);
rejected(&process);
let recovered = joined.recover(RootOutcomeFacts::default()).ok().unwrap();
assert!(
weak.upgrade().is_none(),
"matching join and projection destroy shared state"
);
if stage == 2 {
assert_eq!(recovered.result.unwrap().ok().unwrap(), 41);
} else {
assert!(recovered.primary.is_some());
}
let terminal = recovered.owner.cancel_observed();
assert_eq!(terminal.post_release().cpu_used(), 0);
drop(admit(&process));
process.finish().unwrap();
}
});
}
#[test]
fn derived_view_does_not_retain_obsolete_local_history() {
let process = process();
let (permit, root) = ReservedRequestRoot::try_admitted(
admit(&process),
ContextLabel::checked("app").unwrap(),
)
.ok()
.unwrap();
let view = root.view(RequestViewPhase::Reading).unwrap();
let weak = Arc::downgrade(&view.0);
let derived = view.with_phase(RequestViewPhase::Reading).unwrap();
drop((view, root));
permit.into_execution().cancel_observed();
assert!(
weak.upgrade().is_none(),
"ordinary local history is not retained"
);
rejected(&process);
drop(derived);
assert!(weak.upgrade().is_none());
drop(admit(&process));
process.finish().unwrap();
}
#[test]
fn whole_reservation_shortage_returns_uncalled_factory_and_original_permit() {
let process = process();
let called = Arc::new(std::sync::atomic::AtomicBool::new(false));
let observer = called.clone();
let (original, _make, error) = try_prepare_execution_task(
admit(&process),
ContextLabel::checked("app").unwrap(),
None,
Layout::from_size_align(68152, 8).unwrap(),
move |_, _| {
observer.store(true, std::sync::atomic::Ordering::SeqCst);
Box::pin(std::future::ready(Ok(41u32)))
},
)
.err()
.unwrap();
assert!(matches!(error, ReservedContextError::Storage(_)));
assert!(!called.load(std::sync::atomic::Ordering::SeqCst));
rejected(&process);
drop(original);
drop(admit(&process));
process.finish().unwrap();
}
#[test]
fn child_guard_and_embedded_permission_layouts_are_exact() {
assert_eq!(
Layout::new::<Retained<RootStorage>>(),
Layout::new::<Arc<RootStorage>>()
);
println!(
"runtime_tls_frame={} align={} previous_frame=per_synchronous_nested_poll source_stack_domain=separate",
crate::diagnostics::frame_layout().size(),
crate::diagnostics::frame_layout().align()
);
let [(_, core_root), (_, core_view), (_, core_child)] =
saddle_core::request_context::request_context_layouts();
assert_eq!(
root_demand().unwrap().bytes(),
super::super::storage::shared(Layout::new::<RootStorage>())
.unwrap()
.allocation
.size()
+ core_root.size()
);
assert_eq!(
view_demand(&[(core_view, 1)]).unwrap().bytes(),
super::super::storage::shared(Layout::new::<ViewStorage>())
.unwrap()
.allocation
.size()
+ core_view.size()
);
let child_demand = StorageDemand::embedded(
Layout::new::<[std::sync::atomic::AtomicUsize; 2]>(),
Layout::new::<()>(),
&[(core_child, 1)],
)
.unwrap();
assert_eq!(
child_demand.bytes(),
super::super::storage::shared(Layout::new::<ChildStorage>())
.unwrap()
.allocation
.size()
+ core_child.size()
);
println!(
"root={} view={} child={} permission={} task_status={} context={} public_failure={}",
root_demand().unwrap().bytes(),
view_demand(&[(core_view, 1)]).unwrap().bytes(),
child_demand.bytes(),
Layout::new::<StoragePermit>().size(),
Layout::new::<TaskStatusStorage>().size(),
Layout::new::<ReservedTaskContext>().size(),
Layout::new::<PublicRequestFailure>().size()
);
let process = process();
let (permit, root) = ReservedRequestRoot::try_admitted(
admit(&process),
ContextLabel::checked("app").unwrap(),
)
.ok()
.unwrap();
let call = |rpc: &str, span: u64| {
saddle_core::CallContext::new(
"app".into(),
"module".into(),
"service".into(),
"operation".into(),
saddle_core::TraceId::from_u128(1),
saddle_core::SpanId::from_u64(span),
)
.with_rpc_correlation_id(saddle_core::RpcCorrelationId::new(rpc))
};
let group = RequestIdentityGroup::from_validated(
&call("0", 1),
"request",
"route",
1,
ContextFact::NotEstablished,
)
.unwrap();
root.0
.publisher
.lock()
.unwrap()
.publish(group)
.ok()
.unwrap();
let view = root.view(RequestViewPhase::Reading).unwrap();
let child = view.child(&call("0.1", 2), "request", "child", 1).unwrap();
let weak = Arc::downgrade(child.0.child.as_ref().unwrap());
let next = child.with_phase(RequestViewPhase::Reading).unwrap();
drop((child, view, root));
permit.into_execution().cancel_observed();
assert!(weak.upgrade().is_some(), "shared Core child still charged");
rejected(&process);
drop(next);
assert!(weak.upgrade().is_none());
drop(admit(&process));
process.finish().unwrap();
}
#[test]
fn ready_value_survives_cleanup_panic_after_whole_reservation() {
struct ReadyThenPanic;
impl std::future::Future for ReadyThenPanic {
type Output = Result<u32, ReservedRequestFailure>;
fn poll(
self: std::pin::Pin<&mut Self>,
_: &mut std::task::Context<'_>,
) -> std::task::Poll<Self::Output> {
std::task::Poll::Ready(Ok(41))
}
}
impl Drop for ReadyThenPanic {
fn drop(&mut self) {
panic!("original cleanup description");
}
}
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(async {
let process = process();
let (root, future, mut ticket) = try_prepare_execution_task(
admit(&process),
ContextLabel::checked("app").unwrap(),
None,
Layout::new::<ReadyThenPanic>(),
|_, _| Box::pin(ReadyThenPanic),
)
.ok()
.unwrap();
let mut set = tokio::task::JoinSet::new();
ticket.bind(set.spawn(future).id()).unwrap();
let joined = ticket
.complete(set.join_next().await.unwrap())
.ok()
.unwrap();
drop(root);
let recovered = joined.recover(RootOutcomeFacts::default()).ok().unwrap();
assert_eq!(recovered.result.unwrap().ok().unwrap(), 41);
assert!(recovered.primary.is_none());
assert!(recovered.cleanup.is_some());
assert_eq!(
recovered.owner.cancel_observed().post_release().cpu_used(),
0
);
drop(admit(&process));
process.finish().unwrap();
});
}
#[test]
fn foreign_ticket_is_recoverable_and_unsubmitted_owner_withdraws() {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(async {
let process = process();
let make = || {
try_prepare_rejection_task(
&process.verified_profile(),
ContextLabel::checked("app").unwrap(),
None,
Layout::new::<std::future::Ready<Result<u32, ReservedRequestFailure>>>(),
&[],
(),
|_, _| Box::pin(std::future::ready(Ok(41))),
)
.ok()
.unwrap()
};
let (root_a, future_a, mut ticket_a) = make();
let (root_b, future_b, mut ticket_b) = make();
let mut a = tokio::task::JoinSet::new();
let mut b = tokio::task::JoinSet::new();
ticket_a.bind(a.spawn(future_a).id()).unwrap();
ticket_b.bind(b.spawn(future_b).id()).unwrap();
let output_a = a.join_next().await.unwrap();
let (ticket_b, output_a) = ticket_b
.complete(output_a)
.err()
.expect("foreign output returns both");
let joined_a = ticket_a.complete(output_a).ok().unwrap();
let joined_b = ticket_b
.complete(b.join_next().await.unwrap())
.ok()
.unwrap();
drop((root_a, root_b));
assert_eq!(
joined_a
.recover(Default::default())
.ok()
.unwrap()
.result
.unwrap()
.ok()
.unwrap(),
41
);
assert_eq!(
joined_b
.recover(Default::default())
.ok()
.unwrap()
.result
.unwrap()
.ok()
.unwrap(),
41
);
let (root, future, ticket) = make();
let joined = ticket.withdraw(future).ok().unwrap();
drop(root);
let result = joined.recover(Default::default()).ok().unwrap();
assert!(result.primary.is_some());
assert!(result.result.is_none());
process.finish().unwrap();
});
}
#[test]
fn foreign_source_cannot_be_projected_as_current_request_failure() {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(async {
let process = process();
let foreign_root = ReservedRequestRoot::create(
process
.verified_profile()
.try_rejection_storage(root_demand().unwrap())
.unwrap(),
ContextLabel::checked("app").unwrap(),
)
.unwrap();
let foreign_view = foreign_root.view(RequestViewPhase::Reading).unwrap();
let failure = foreign_view.source_error(
&std::io::Error::other("foreign original"),
None,
saddle_core::DiagnosticStage::RequestDb,
RootRequestEvent::Finalization,
Default::default(),
);
let id = failure.occurrence();
let (root, future, mut ticket) = try_prepare_execution_task(
admit(&process),
ContextLabel::checked("app").unwrap(),
None,
Layout::new::<std::future::Ready<Result<u32, ReservedRequestFailure>>>(),
move |_, _| Box::pin(std::future::ready(Err::<u32, _>(failure))),
)
.ok()
.unwrap();
let mut tasks = tokio::task::JoinSet::new();
ticket.bind(tasks.spawn(future).id()).unwrap();
let joined = ticket
.complete(tasks.join_next().await.unwrap())
.ok()
.unwrap();
let mut joined = joined
.recover(Default::default())
.err()
.expect("foreign source fails closed");
let original = joined.take_foreign_failure().unwrap();
assert_eq!(
serde_json::to_value(original.occurrence()).unwrap(),
serde_json::to_value(id).unwrap()
);
let public = original
.finish(&foreign_view, None, Default::default())
.ok()
.unwrap();
let recovered = joined.recover(Default::default()).ok().unwrap();
drop((root, foreign_view, foreign_root));
assert!(recovered.result.is_none());
assert!(recovered.primary.is_some());
assert_eq!(
recovered.owner.cancel_observed().post_release().cpu_used(),
0
);
std::hint::black_box(public);
process.finish().unwrap();
});
}
#[test]
fn panic_cleanup_originals_use_one_root_and_survive_file_readback() {
const CHILD: &str = "R_RESERVED_WRITER_CHILD";
let Some(path) = std::env::var_os(CHILD) else {
let base = std::env::var_os("R_COMPONENT_EVIDENCE")
.map(std::path::PathBuf::from)
.unwrap_or_else(std::env::temp_dir);
let path = base.join(format!("reserved-writer-{}", std::process::id()));
std::fs::create_dir(&path).unwrap();
let mut child=std::process::Command::new(std::env::current_exe().unwrap()).args(["--exact",
"request_task::reserved::tests::panic_cleanup_originals_use_one_root_and_survive_file_readback","--nocapture"])
.env(CHILD,&path).spawn().unwrap();
let start = std::time::Instant::now();
loop {
if let Some(status) = child.try_wait().unwrap() {
assert!(status.success());
break;
}
if start.elapsed() > std::time::Duration::from_secs(15) {
child.kill().unwrap();
child.wait().unwrap();
panic!("owned writer child timeout");
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
println!("writer evidence {}", path.display());
return;
};
let output = saddle_observability::EmergencyDiagnostics::start(
&saddle_observability::FileLoggingConfig::new(
std::path::PathBuf::from(&path),
saddle_observability::Rotation::Daily,
),
)
.unwrap();
std::panic::set_hook(Box::new(|info| {
let _ = crate::diagnostics::capture_current_panic(info);
}));
struct PanicDrop;
impl std::future::Future for PanicDrop {
type Output = Result<u32, ReservedRequestFailure>;
fn poll(
self: std::pin::Pin<&mut Self>,
_: &mut std::task::Context<'_>,
) -> std::task::Poll<Self::Output> {
panic!("R_ORIGINAL_PRIMARY 中文");
}
}
impl Drop for PanicDrop {
fn drop(&mut self) {
panic!("R_ORIGINAL_CLEANUP 中文");
}
}
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap()
.block_on(async {
let process = process();
let (root, future, mut ticket) = try_prepare_execution_task(
admit(&process),
ContextLabel::checked("app").unwrap(),
Some(output.handle()),
Layout::new::<PanicDrop>(),
|_, mut context| {
let call = saddle_core::CallContext::new(
"app".into(),
"module".into(),
"svc".into(),
"operation".into(),
saddle_core::TraceId::from_u128(1),
saddle_core::SpanId::from_u64(1),
)
.with_rpc_correlation_id(saddle_core::RpcCorrelationId::new("0"));
context
.publish(
RequestIdentityGroup::from_validated(
&call,
"request",
"route",
1,
ContextFact::NotEstablished,
)
.unwrap(),
)
.unwrap();
Box::pin(PanicDrop)
},
)
.ok()
.unwrap();
let mut tasks = tokio::task::JoinSet::new();
ticket.bind(tasks.spawn(future).id()).unwrap();
let joined = ticket
.complete(tasks.join_next().await.unwrap())
.ok()
.unwrap();
let primary =
serde_json::to_value(joined.primary.as_ref().unwrap().occurrence()).unwrap();
let cleanup =
serde_json::to_value(joined.cleanup.as_ref().unwrap().occurrence()).unwrap();
assert_eq!(cleanup["primary_diagnostic_id"], primary["diagnostic_id"]);
let recovered = joined.recover(Default::default()).ok().unwrap();
drop(root);
assert_eq!(
recovered.owner.cancel_observed().post_release().cpu_used(),
0
);
drop(admit(&process));
process.finish().unwrap();
});
let exit = crate::diagnostics::close_output(
output,
Some(std::time::Instant::now() + std::time::Duration::from_secs(2)),
);
assert_eq!(exit.snapshot.enqueued, exit.snapshot.written);
assert_eq!(exit.snapshot.dropped, 0);
let raw =
std::fs::read_to_string(std::path::PathBuf::from(path).join("saddle.emergency.log"))
.unwrap();
let rows: Vec<serde_json::Value> = raw
.lines()
.map(|line| serde_json::from_str(line).unwrap())
.collect();
let descriptions: Vec<_> = rows
.iter()
.filter(|row| row["channel"] == "description")
.collect();
assert_eq!(descriptions.len(), 2);
assert_eq!(descriptions[0]["payload"], "R_ORIGINAL_PRIMARY 中文");
assert_eq!(descriptions[1]["payload"], "R_ORIGINAL_CLEANUP 中文");
let headers: Vec<_> = rows
.iter()
.filter(|row| row["channel"] == "context")
.map(|row| {
serde_json::from_str::<serde_json::Value>(row["payload"].as_str().unwrap()).unwrap()
})
.collect();
assert_eq!(headers.len(), 2);
for header in &headers {
assert_eq!(header["context"]["request"]["value"], "request");
assert_eq!(header["context"]["task"]["state"], "present");
assert_eq!(header["facts"]["stack_status"], "unavailable_deferred");
assert!(header["facts"]["origin"]["line"].as_u64().unwrap() > 0);
}
assert_eq!(headers[0]["context"], headers[1]["context"]);
}
}