use super::*;
use saddle_boundary::reserved_diagnostics::{ReservedDelivery, ReservedDeliveryOutcome};
use saddle_observability::root_diagnostic::{RootOutcomeFacts, RootRequestEvent};
use saddle_observability::{AdmissionCapacitySubmission, AdmissionConstruction};
use saddle_runtime::profusegw::ReservedDispatchOwner;
use saddle_runtime::request_task::reserved::{
ReservedBorrowedFuture, ReservedRequestFailure, ReservedRequestView, ReservedTaskContext,
};
enum UnrootedEntryError<'a> {
Admission(&'a saddle_admission::AdmissionError),
Admitted(&'a saddle_runtime::profusegw::ReservedDispatchConstructionFailure),
Rejection(&'a saddle_runtime::request_task::reserved::ReservedContextError),
}
#[track_caller]
fn record_unrooted_failure(
application: &saddle_core::ContextLabel,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
error: UnrootedEntryError<'_>,
) -> saddle_observability::root_diagnostic::UnrootedCaptureFacts {
use saddle_observability::root_diagnostic::UnrootedDiagnosticScope;
use saddle_runtime::profusegw::ReservedDispatchConstructionFailure as F;
struct DebugOnly<'a, T>(&'a T);
impl<T: std::fmt::Debug> std::fmt::Debug for DebugOnly<'_, T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.0.fmt(f)
}
}
impl<T: std::fmt::Debug> std::fmt::Display for DebugOnly<'_, T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
std::fmt::Debug::fmt(self.0, f)
}
}
let scope =
UnrootedDiagnosticScope::new(application, saddle_core::RequestViewPhase::Admitted, output);
let axes = RootOutcomeFacts {
axes: saddle_core::DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Failed,
..Default::default()
},
..Default::default()
};
match error {
UnrootedEntryError::Admission(error)
| UnrootedEntryError::Admitted(F::Storage(error) | F::RequestIdentity(error)) => scope
.source_error(
error,
saddle_core::DiagnosticStage::RequestAdmission,
RootRequestEvent::Admission,
axes,
),
UnrootedEntryError::Admitted(F::Deadline(error)) => scope.source_description(
&DebugOnly(error),
saddle_core::DiagnosticStage::RequestAdmission,
RootRequestEvent::Admission,
axes,
),
UnrootedEntryError::Admitted(F::Context(error)) | UnrootedEntryError::Rejection(error) => {
scope.source_description(
&DebugOnly(error),
saddle_core::DiagnosticStage::RequestAdmission,
RootRequestEvent::Admission,
axes,
)
}
}
}
#[cfg(test)]
mod unrooted_tests {
use super::*;
#[test]
fn three_formal_capture_points_preserve_originals_without_a_root() {
use saddle_observability::{
DiagnosticSubmission, EmergencyDiagnostics, FileLoggingConfig, Rotation,
};
use saddle_runtime::{
profusegw::ReservedDispatchConstructionFailure as F,
request_task::reserved::ReservedContextError,
};
let directory =
std::env::temp_dir().join(format!("saddle-s-unrooted-{}", std::process::id()));
std::fs::create_dir_all(&directory).unwrap();
let mut writer = EmergencyDiagnostics::start(&FileLoggingConfig::new(
directory.clone(),
Rotation::Daily,
))
.unwrap();
let output = writer.handle();
let application = saddle_core::ContextLabel::checked("formal-s").unwrap();
let admission = saddle_admission::AdmissionError::AccountClosed;
let admitted = F::Storage(saddle_admission::AdmissionError::SizeOverflow);
let rejected =
ReservedContextError::Storage(saddle_admission::AdmissionError::NoAccountSlot);
let mut expected = Vec::new();
for (error, text) in [
(
UnrootedEntryError::Admission(&admission),
format!("{admission:?}"),
),
(
UnrootedEntryError::Admitted(&admitted),
format!("{:?}", saddle_admission::AdmissionError::SizeOverflow),
),
(
UnrootedEntryError::Rejection(&rejected),
format!("{rejected:?}"),
),
] {
let facts = record_unrooted_failure(&application, Some(&output), error);
assert_eq!(facts.submission(), DiagnosticSubmission::Enqueued);
assert_eq!(
facts.original_capture(),
saddle_observability::root_diagnostic::OriginalCaptureState::CompleteEnqueued
);
expected.push((serde_json::to_value(facts.occurrence()).unwrap(), text));
}
let unavailable = record_unrooted_failure(
&application,
None,
UnrootedEntryError::Admission(&admission),
);
assert_eq!(
unavailable.submission(),
DiagnosticSubmission::OutputUnavailable
);
let _ = writer.shutdown();
let contents = std::fs::read_to_string(directory.join("saddle.emergency.log")).unwrap();
let rows: Vec<serde_json::Value> = contents
.lines()
.map(|line| serde_json::from_str(line).unwrap())
.collect();
for (id, text) in expected {
let records: Vec<_> = rows.iter().filter(|row| row["occurrence"] == id).collect();
let field = |channel: &str| {
records
.iter()
.filter(|row| row["channel"] == channel)
.map(|row| row["payload"].as_str().unwrap())
.collect::<String>()
};
assert_eq!(field("debug"), text);
assert_eq!(field("description"), text);
let header: serde_json::Value = serde_json::from_str(&field("context")).unwrap();
assert_eq!(header["context"]["application"]["value"], "formal-s");
assert_eq!(header["context"]["lifecycle"]["value"], "admitted");
for key in [
"task",
"scope",
"trace_id",
"rpc_id",
"local_request",
"zone",
"operation",
] {
assert_eq!(header["context"][key]["state"], "not_established");
}
assert!(records.iter().any(|row| row["channel"] == "terminal"));
}
println!("S_UNROOTED_THREE_POINTS PASS raw={}", directory.display());
}
}
pub(super) struct EntryOwner {
pub database: Option<crate::database_capability::DatabaseRequestCompletion>,
pub view: Option<ReservedRequestView>,
pub delivery: Option<ReservedDeliveryOutcome>,
pub admission: Option<AdmissionCapacitySubmission>,
}
impl EntryOwner {
fn delivery_facts(&self) -> RootOutcomeFacts {
let (delivery, bytes_written) = match &self.delivery {
Some(ReservedDeliveryOutcome::Complete(axes))
| Some(ReservedDeliveryOutcome::Failed { axes, .. }) => {
(axes.delivery, axes.bytes_written)
}
None => (saddle_core::ResponseDelivery::NotStarted, Some(0)),
};
RootOutcomeFacts {
axes: saddle_core::DiagnosticOutcomeAxes {
delivery,
bytes_written,
..Default::default()
},
..Default::default()
}
}
fn new() -> Self {
Self {
database: None,
view: None,
delivery: None,
admission: None,
}
}
async fn cleanup(
&mut self,
observer: &saddle_observability::Observer,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
) {
if self.delivery.is_none() {
if let Some(view) = &self.view {
let _ = view.ordinary(
observer,
RootRequestEvent::Response,
RootOutcomeFacts {
axes: saddle_core::DiagnosticOutcomeAxes {
delivery: saddle_core::ResponseDelivery::NotStarted,
bytes_written: Some(0),
..Default::default()
},
..Default::default()
},
);
}
}
if let Some(database) = &self.database {
database.finish(true).await;
for failure in [
database.take_reserved_failure(),
database.take_encoding_failure(),
]
.into_iter()
.flatten()
{
if failure
.finish(
self.view.as_ref().expect("original DB view"),
output,
Default::default(),
)
.is_err()
{
unreachable!("original request owns DB/encoding receipt");
}
}
}
if let Some(outcome) = self.delivery.take() {
let view = self.view.as_ref().expect("delivery requires actual view");
match outcome {
ReservedDeliveryOutcome::Complete(axes) => {
let _ = view.ordinary(
observer,
RootRequestEvent::Response,
RootOutcomeFacts {
axes,
..Default::default()
},
);
}
ReservedDeliveryOutcome::Failed { axes, failure } => {
let (_, receipt) = failure.into_parts();
if receipt
.finish(
view,
output,
RootOutcomeFacts {
axes,
..Default::default()
},
)
.is_err()
{
unreachable!("delivery is owned by the original request");
}
}
}
}
}
}
pub(super) struct NormalInput<B, C, Dispatch> {
pub factory: Option<consumer_storage::NormalFactory<B, C, Dispatch>>,
pub retained: EntryOwner,
}
struct RejectedInput {
socket: Option<TcpStream>,
code: u16,
admission: Option<saddle_runtime::profusegw::ProfuseGwAdmissionEvent>,
retained: EntryOwner,
observer: saddle_observability::Observer,
output: Option<saddle_observability::EmergencyDiagnosticHandle>,
}
async fn rejected_body(
owner: &mut RejectedInput,
context: ReservedTaskContext,
) -> std::result::Result<(), ReservedRequestFailure> {
let view = context.view();
owner.retained.view = Some(view.clone());
if let Some(admission) = owner.admission.take() {
owner.retained.admission = Some(view.record_admission(
admission,
&owner.observer,
owner.output.as_ref(),
AdmissionConstruction::Ready,
));
}
let mut socket = owner.socket.take().expect("original rejected socket");
let deadline = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.ok()
.and_then(|now| i64::try_from(now.as_millis()).ok())
.unwrap_or(0)
.saturating_add(50);
deliver(
&mut socket,
&status(owner.code),
deadline,
view,
owner.output.as_ref(),
&mut owner.retained.delivery,
)
.await;
Ok(())
}
fn rejected_factory(
owner: &mut RejectedInput,
context: ReservedTaskContext,
) -> ReservedBorrowedFuture<'_, ()> {
Box::pin(rejected_body(owner, context))
}
fn rejected_layout() -> std::alloc::Layout {
fn out<A, R>(_: impl FnOnce(A) -> R) -> std::alloc::Layout {
std::alloc::Layout::new::<R>()
}
out(
|(owner, context): (&'static mut RejectedInput, ReservedTaskContext)| {
rejected_body(owner, context)
},
)
}
async fn finish_normal<B, C, Dispatch>(
mut joined: saddle_runtime::request_task::reserved::ReservedTaskJoined<
Owner<B, C, Dispatch>,
(),
>,
observer: &saddle_observability::Observer,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
) {
if let Some(admission) = joined.owner_mut().2.take() {
let submission =
joined.record_admission(admission, observer, AdmissionConstruction::Cancelled);
joined.owner_mut().1.retained.admission = Some(submission);
}
let owner = joined.owner_mut();
let terminal_facts = owner.1.retained.delivery_facts();
owner.1.retained.cleanup(observer, output).await;
if let Some(dispatch) = owner.0.take() {
dispatch.cancel();
}
let recovered = joined
.recover(terminal_facts)
.unwrap_or_else(|_| unreachable!("original request failure"));
drop(recovered);
}
async fn finish_rejected(
mut joined: saddle_runtime::request_task::reserved::ReservedTaskJoined<RejectedInput, ()>,
) {
if let Some(admission) = joined.owner_mut().admission.take() {
let observer = joined.owner_mut().observer.clone();
let submission =
joined.record_admission(admission, &observer, AdmissionConstruction::Cancelled);
joined.owner_mut().retained.admission = Some(submission);
}
let owner = joined.owner_mut();
let terminal_facts = owner.retained.delivery_facts();
owner
.retained
.cleanup(&owner.observer, owner.output.as_ref())
.await;
drop(
joined
.recover(terminal_facts)
.unwrap_or_else(|_| unreachable!("original rejection failure")),
);
}
pub(super) struct Prepared<B, C, Dispatch> {
application: saddle_core::ContextLabel,
normal: saddle_runtime::request_task::reserved_set::ReservedTaskCollection<
Owner<B, C, Dispatch>,
(),
>,
rejected: saddle_runtime::request_task::reserved_set::ReservedTaskCollection<RejectedInput, ()>,
indirect: [(std::alloc::Layout, usize); 6],
}
pub(super) fn prepare<B, C, Dispatch, DispatchFuture>(
application: &str,
lease: &saddle_runtime::profusegw::ProfuseGwProcessLease,
_: &Dispatch,
) -> std::result::Result<Prepared<B, C, Dispatch>, saddle_admission::AdmissionError>
where
B: Send + Sync + 'static,
C: Clone + Send + Sync + 'static,
Dispatch: Fn(
saddle_boundary::ingress::AcceptedIngress,
C,
BusinessConfig<B>,
crate::database_capability::DatabaseRequest,
) -> DispatchFuture
+ Send
+ Sync
+ 'static,
DispatchFuture: Future<Output = Result<Vec<u8>>> + Send + 'static,
{
let application = saddle_core::ContextLabel::checked(application)
.map_err(|_| saddle_admission::AdmissionError::InvalidConfiguration)?;
let storage = lease.try_reserved_task_storage()?;
let normal = storage.normal::<Owner<B, C, Dispatch>, ()>()?;
let rejected = storage.rejection::<RejectedInput, ()>()?;
let database = crate::database_capability::request_storage_layouts()
.map_err(|_| saddle_admission::AdmissionError::SizeOverflow)?;
let context = consumer_storage::shared_layout::<ReservedTaskContext>()
.map_err(|_| saddle_admission::AdmissionError::SizeOverflow)?;
let call_counter = consumer_storage::shared_layout::<std::sync::atomic::AtomicU64>()
.map_err(|_| saddle_admission::AdmissionError::SizeOverflow)?;
Ok(Prepared {
application,
normal,
rejected,
indirect: [
(database.state.allocation, 1),
(database.cancellation_notify.allocation, 1),
(database.outbound_retentions.allocation, 1),
(context.allocation, 1),
(
std::alloc::Layout::new::<tokio::sync::futures::OwnedNotified>(),
1,
),
(call_counter.allocation, 1),
],
})
}
pub(super) async fn run<B, C, Dispatch, DispatchFuture>(
prepared: Prepared<B, C, Dispatch>,
listener: TcpListener,
adapter: saddle_boundary::ingress::ProfuseGwListenerAdapter,
deployment: C,
dispatch: Dispatch,
business: BusinessConfig<B>,
ingress_token: Option<Arc<IngressToken>>,
lease: saddle_runtime::profusegw::ProfuseGwProcessLease,
observer: saddle_observability::Observer,
database: Option<Arc<ManagedDatabaseState>>,
output: Option<saddle_observability::EmergencyDiagnosticHandle>,
stop: Arc<AtomicBool>,
) where
B: Send + Sync + 'static,
C: Clone + Send + Sync + 'static,
Dispatch: Fn(
saddle_boundary::ingress::AcceptedIngress,
C,
BusinessConfig<B>,
crate::database_capability::DatabaseRequest,
) -> DispatchFuture
+ Send
+ Sync
+ 'static,
DispatchFuture: Future<Output = Result<Vec<u8>>> + Send + 'static,
{
use saddle_observability::{
AdmissionConstruction as Construction, AdmissionConstructionFailure as Failure,
};
use saddle_runtime::profusegw::{
ReservedDispatchConstructionFailure, ReservedDispatchOutcome, ReservedDispatchRejection,
};
use saddle_runtime::request_task::reserved_set::ReservedCollectionJoin;
let Prepared {
application,
mut normal,
mut rejected,
indirect,
} = prepared;
let dispatch = Arc::new(dispatch);
while !stop.load(Ordering::Acquire) {
tokio::select! {
accepted=listener.accept(),if rejected.len()<16 => {
let (socket,_)=match accepted {Ok(value)=>value,Err(error)=>{
let diagnostic=saddle_boundary::diagnostics::http_io(saddle_boundary::diagnostics::HttpIoStage::Listener,&error);
let _=diagnostics::record(&diagnostic,None);
continue;
}};
let input=NormalInput{factory:Some(consumer_storage::NormalFactory{
socket,adapter:adapter.clone(),deployment:Some(deployment.clone()),business:Some(business.clone()),
ingress_token:ingress_token.clone(),dispatch:dispatch.clone(),observer:observer.clone(),admission:None,
database:database.clone(),diagnostic_handle:output.clone(),
}),retained:EntryOwner::new()};
match lease.try_reserved_dispatch(application.clone(),output.clone(),body_layout::<B,C,Dispatch,DispatchFuture>(),&indirect,input,factory::<B,C,Dispatch,DispatchFuture>) {
ReservedDispatchOutcome::Ready{root,future,ticket}=> {
if let Err((future,ticket,_error))=normal.spawn(future,ticket) {
let joined=ticket.withdraw(future).unwrap_or_else(|_|unreachable!("original unspawned pair"));
finish_normal(joined,&observer,output.as_ref()).await;
}
drop(root);
}
ReservedDispatchOutcome::Rejected{mut input,make:_,reason}=> {
let (code,admission)=match reason {
ReservedDispatchRejection::Capacity(admission)=>(503,Some(admission)),
ReservedDispatchRejection::Admitted{admission,failure}=> {
let _original=record_unrooted_failure(&application,output.as_ref(),UnrootedEntryError::Admitted(&failure));
let failure=match failure {
ReservedDispatchConstructionFailure::Storage(_)=>Failure::Storage,
ReservedDispatchConstructionFailure::Deadline(_)=>Failure::Deadline,
ReservedDispatchConstructionFailure::RequestIdentity(_)=>Failure::RequestIdentity,
ReservedDispatchConstructionFailure::Context(_)=>Failure::Context,
};
let _submissions=admission.record_unrooted(&observer,output.as_ref(),&application,saddle_core::RequestViewPhase::Admitted,Construction::Failed(failure));
(500,None)
}
ReservedDispatchRejection::Admission(error)=> {
let _original=record_unrooted_failure(&application,output.as_ref(),UnrootedEntryError::Admission(&error));
(500,None)
},
};
let factory=input.factory.take().expect("failed construction returns uncalled factory");
let input=RejectedInput{socket:Some(factory.socket),code,admission,retained:EntryOwner::new(),observer:observer.clone(),output:output.clone()};
match lease.try_reserved_rejection(application.clone(),output.clone(),rejected_layout(),&[],input,rejected_factory) {
Ok((root,future,ticket))=> {
if let Err((future,ticket,_error))=rejected.spawn(future,ticket) {
let joined=ticket.withdraw(future).unwrap_or_else(|_|unreachable!("original rejected pair"));
finish_rejected(joined).await;
}
drop(root);
}
Err((mut input,_make,error))=> {
let _original=record_unrooted_failure(&application,output.as_ref(),UnrootedEntryError::Rejection(&error));
if let Some(admission)=input.admission.take() {
let failure=match error {
saddle_runtime::request_task::reserved::ReservedContextError::Storage(saddle_admission::AdmissionError::CapacityRejected|saddle_admission::AdmissionError::NoAccountSlot)=>Failure::RejectionTaskSlot,
saddle_runtime::request_task::reserved::ReservedContextError::Storage(_)=>Failure::RejectionTaskBytes,
saddle_runtime::request_task::reserved::ReservedContextError::Context(_)=>Failure::Context,
};
input.retained.admission=Some(admission.record_unrooted(&observer,output.as_ref(),&application,saddle_core::RequestViewPhase::Admitted,Construction::Failed(failure)));
}
drop(input);
}
}
}
}
}
Some(join)=normal.join_next(),if !normal.is_empty()=> {
let ReservedCollectionJoin::Matched(joined)=join else {unreachable!("collection retains its matching ticket")};
finish_normal(joined,&observer,output.as_ref()).await;
}
Some(join)=rejected.join_next(),if !rejected.is_empty()=> {
let ReservedCollectionJoin::Matched(joined)=join else {unreachable!("collection retains its matching rejection ticket")};
finish_rejected(joined).await;
}
()=tokio::time::sleep(Duration::from_millis(1))=>{},
}
}
while let Some(join) = normal.join_next().await {
let ReservedCollectionJoin::Matched(joined) = join else {
unreachable!("matching shutdown join")
};
finish_normal(joined, &observer, output.as_ref()).await;
}
while let Some(join) = rejected.join_next().await {
let ReservedCollectionJoin::Matched(joined) = join else {
unreachable!("matching shutdown rejection")
};
finish_rejected(joined).await;
}
drop((normal, rejected));
}
type Owner<B, C, Dispatch> = ReservedDispatchOwner<NormalInput<B, C, Dispatch>>;
struct WriteOwner<'slot, 'output> {
progress: Option<ReservedDelivery<'output>>,
retained: &'slot mut Option<ReservedDeliveryOutcome>,
}
impl Drop for WriteOwner<'_, '_> {
fn drop(&mut self) {
if let Some(mut progress) = self.progress.take() {
progress.cancelled();
*self.retained = Some(progress.finish());
}
}
}
async fn deliver<W: tokio::io::AsyncWrite + Unpin>(
socket: &mut W,
bytes: &[u8],
deadline: i64,
view: ReservedRequestView,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
retained: &mut Option<ReservedDeliveryOutcome>,
) {
let mut owner = WriteOwner {
progress: Some(ReservedDelivery::new(view, output)),
retained,
};
{
let progress = owner.progress.as_mut().unwrap();
tokio::select! {
biased;
() = sleep_until_unix_ms(deadline) => progress.timed_out(),
() = async {
let mut offset = 0;
while offset < bytes.len() && progress.begin_write() {
let result = socket.write(&bytes[offset..]).await;
let count = result.as_ref().copied().unwrap_or(0);
if !matches!(progress.record_write(result), saddle_boundary::request_diagnostics::WriteStep::Progress) { return; }
offset += count;
}
if offset == bytes.len() {
progress.local_write_complete();
if let Err(error) = socket.shutdown().await { progress.io_failure(&error); }
}
} => {},
}
}
*owner.retained = Some(owner.progress.take().unwrap().finish());
}
#[cfg(test)]
#[path = "reserved_entry_tests.rs"]
mod lifecycle_tests;
impl<B, C, Dispatch> consumer_storage::NormalFactory<B, C, Dispatch> {
pub(super) async fn reserved_body<DispatchFuture>(
&mut self,
owner: (
&mut Option<saddle_runtime::profusegw::ProfuseGwManagedDispatch>,
&mut EntryOwner,
&mut Option<saddle_runtime::profusegw::ProfuseGwAdmissionEvent>,
),
mut task: ReservedTaskContext,
) -> std::result::Result<(), ReservedRequestFailure>
where
B: Send + Sync + 'static,
C: Send + 'static,
Dispatch: Fn(
saddle_boundary::ingress::AcceptedIngress,
C,
BusinessConfig<B>,
crate::database_capability::DatabaseRequest,
) -> DispatchFuture
+ Send
+ Sync
+ 'static,
DispatchFuture: Future<Output = Result<Vec<u8>>> + Send + 'static,
{
let prepared = self.reserved_ingress(owner.0, owner.2, &mut task).await;
owner.1.view = Some(task.view());
let prepared = match prepared {
Ok(prepared) => prepared,
Err(failure) => {
if let Some(admission) = owner.2.take() {
owner.1.admission = Some(task.view().record_admission(
admission,
&self.observer,
self.diagnostic_handle.as_ref(),
AdmissionConstruction::Failed(
saddle_observability::AdmissionConstructionFailure::Context,
),
));
}
let deadline = owner.0.as_ref().unwrap().deadline_unix_ms();
deliver(
&mut self.socket,
&status(failure.status),
deadline,
task.view(),
self.diagnostic_handle.as_ref(),
&mut owner.1.delivery,
)
.await;
return Err(failure.source);
}
};
owner.1.admission = Some(prepared.admission_submission);
let handler_view = task
.view()
.with_phase(saddle_core::request_context::RequestViewPhase::Handler)
.map_err(|error| {
#[derive(Debug)]
struct Description(saddle_runtime::request_task::reserved::ReservedContextError);
impl std::fmt::Display for Description {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{:?}", self.0)
}
}
let code = saddle_core::DiagnosticCode::new("service.handler.context").unwrap();
task.view().source_description(
&Description(error),
saddle_core::BoundedDiagnostic::capture(
saddle_core::DiagnosticCategory::ExpectedRejection,
saddle_core::CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
saddle_core::DiagnosticStage::RequestHandler,
code,
),
),
code,
self.diagnostic_handle.as_ref(),
RootRequestEvent::Handler,
Default::default(),
)
})?;
task.observe_view(handler_view)
.unwrap_or_else(|_| unreachable!("same task root"));
owner.1.view = Some(task.view());
let task = Arc::new(task);
let view = task.view();
let deadline = owner.0.as_ref().unwrap().deadline_unix_ms();
let process = self
.database
.as_ref()
.and_then(|db| db.capability.lock().unwrap().as_ref().cloned());
let legacy = saddle_observability::RequestDiagnosticScope::output_unavailable(
prepared.call.context(),
&prepared.event,
);
let (request, completion) = crate::database_capability::DatabaseRequest::new(
process,
owner.0.take().unwrap(),
prepared.call.context().clone(),
prepared.event,
self.diagnostic_handle.clone(),
None,
&prepared.accepted.zone,
legacy,
Some(task),
Some(self.observer.clone()),
);
owner.1.database = Some(completion);
enum Terminal {
Response(Result<Vec<u8>>),
Deadline,
Disconnected,
}
let terminal = {
let dispatch = (self.dispatch)(
prepared.accepted,
self.deployment.take().expect("single deployment input"),
self.business.take().expect("single business input"),
request,
);
tokio::pin!(dispatch);
tokio::select! {
biased;
()=sleep_until_unix_ms(deadline)=>Terminal::Deadline,
()=client_disconnected(&self.socket)=>Terminal::Disconnected,
result=&mut dispatch=>Terminal::Response(result),
}
};
match terminal {
Terminal::Response(Ok(bytes)) => {
deliver(
&mut self.socket,
&bytes,
deadline,
view,
self.diagnostic_handle.as_ref(),
&mut owner.1.delivery,
)
.await;
prepared.call.succeed();
owner.1.database.as_ref().unwrap().finish(false).await;
Ok(())
}
Terminal::Response(Err(error)) => {
if let Some(operation) = owner
.1
.database
.as_ref()
.unwrap()
.take_encoding_supervision()
{
if matches!(operation, saddle_core::OperationOutcome::TimedOut) {
deliver_deadline_now(
&self.socket,
view,
self.diagnostic_handle.as_ref(),
&mut owner.1.delivery,
);
}
owner.1.database.as_ref().unwrap().finish(true).await;
drop(prepared.call);
return Ok(());
}
let code = saddle_core::DiagnosticCode::new("service.dispatch.failure").unwrap();
let source = owner
.1
.database
.as_ref()
.unwrap()
.take_encoding_failure()
.unwrap_or_else(|| {
view.source_error_with_facts(
&error,
saddle_core::BoundedDiagnostic::capture(
saddle_core::DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
saddle_core::DiagnosticStage::RequestHandler,
code,
),
),
code,
self.diagnostic_handle.as_ref(),
RootRequestEvent::Handler,
RootOutcomeFacts::default(),
)
});
deliver(
&mut self.socket,
&status(500),
deadline,
view,
self.diagnostic_handle.as_ref(),
&mut owner.1.delivery,
)
.await;
prepared.call.fail(&error);
owner.1.database.as_ref().unwrap().finish(false).await;
Err(source)
}
stopped @ (Terminal::Deadline | Terminal::Disconnected) => {
let timed_out = matches!(stopped, Terminal::Deadline);
let operation = if timed_out {
saddle_core::OperationOutcome::TimedOut
} else {
saddle_core::OperationOutcome::Cancelled
};
let code = saddle_core::DiagnosticCode::new(if timed_out {
"service.request.deadline"
} else {
"service.request.disconnected"
})
.unwrap();
let source = view.source_description(
&if timed_out {
"original request deadline elapsed"
} else {
"original peer connection closed"
},
saddle_core::BoundedDiagnostic::capture(
saddle_core::DiagnosticCategory::ExpectedRejection,
saddle_core::CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
saddle_core::DiagnosticStage::RequestHandler,
code,
),
),
code,
self.diagnostic_handle.as_ref(),
RootRequestEvent::Handler,
RootOutcomeFacts {
axes: saddle_core::DiagnosticOutcomeAxes {
operation,
..Default::default()
},
..Default::default()
},
);
owner
.1
.database
.as_ref()
.unwrap()
.finish_handler(operation, None);
if timed_out {
deliver_deadline_now(
&self.socket,
view.clone(),
self.diagnostic_handle.as_ref(),
&mut owner.1.delivery,
);
}
owner.1.database.as_ref().unwrap().finish(true).await;
drop(prepared.call);
Err(source)
}
}
}
}
fn deliver_deadline_now(
socket: &TcpStream,
view: ReservedRequestView,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
retained: &mut Option<ReservedDeliveryOutcome>,
) {
let bytes = deadline_terminal_bytes();
let mut progress = ReservedDelivery::new(view, output);
let mut offset = 0;
while offset < bytes.len() {
let result = socket.try_write(&bytes[offset..]);
let count = result.as_ref().copied().unwrap_or(0);
if matches!(
progress.record_write(result),
saddle_boundary::request_diagnostics::WriteStep::Stopped
) {
break;
}
offset += count;
}
if offset == bytes.len() {
progress.local_write_complete();
}
*retained = Some(progress.finish());
}
pub(super) fn factory<B, C, Dispatch, DispatchFuture>(
owner: &mut Owner<B, C, Dispatch>,
context: ReservedTaskContext,
) -> ReservedBorrowedFuture<'_, ()>
where
B: Send + Sync + 'static,
C: Send + 'static,
Dispatch: Fn(
saddle_boundary::ingress::AcceptedIngress,
C,
BusinessConfig<B>,
crate::database_capability::DatabaseRequest,
) -> DispatchFuture
+ Send
+ Sync
+ 'static,
DispatchFuture: Future<Output = Result<Vec<u8>>> + Send + 'static,
{
Box::pin(normal_body::<B, C, Dispatch, DispatchFuture>(
owner, context,
))
}
pub(super) fn body_layout<B, C, Dispatch, DispatchFuture>() -> std::alloc::Layout
where
B: Send + Sync + 'static,
C: Send + 'static,
Dispatch: Fn(
saddle_boundary::ingress::AcceptedIngress,
C,
BusinessConfig<B>,
crate::database_capability::DatabaseRequest,
) -> DispatchFuture
+ Send
+ Sync
+ 'static,
DispatchFuture: Future<Output = Result<Vec<u8>>> + Send + 'static,
{
fn output<A, R>(_: impl FnOnce(A) -> R) -> std::alloc::Layout {
std::alloc::Layout::new::<R>()
}
output(
|(owner, context): (&'static mut Owner<B, C, Dispatch>, ReservedTaskContext)| {
normal_body::<B, C, Dispatch, DispatchFuture>(owner, context)
},
)
}
async fn normal_body<B, C, Dispatch, DispatchFuture>(
owner: &mut Owner<B, C, Dispatch>,
context: ReservedTaskContext,
) -> std::result::Result<(), ReservedRequestFailure>
where
B: Send + Sync + 'static,
C: Send + 'static,
Dispatch: Fn(
saddle_boundary::ingress::AcceptedIngress,
C,
BusinessConfig<B>,
crate::database_capability::DatabaseRequest,
) -> DispatchFuture
+ Send
+ Sync
+ 'static,
DispatchFuture: Future<Output = Result<Vec<u8>>> + Send + 'static,
{
let factory = owner.1.factory.as_mut().expect("single original factory");
factory
.reserved_body((&mut owner.0, &mut owner.1.retained, &mut owner.2), context)
.await
}
#[doc(hidden)]
pub fn layouts_for<B, C, Dispatch, DispatchFuture>(
_: &Dispatch,
) -> std::result::Result<
(std::alloc::Layout, std::alloc::Layout, usize, [usize; 5]),
saddle_admission::AdmissionError,
>
where
B: Send + Sync + 'static,
C: Send + 'static,
Dispatch: Fn(
saddle_boundary::ingress::AcceptedIngress,
C,
BusinessConfig<B>,
crate::database_capability::DatabaseRequest,
) -> DispatchFuture
+ Send
+ Sync
+ 'static,
DispatchFuture: Future<Output = Result<Vec<u8>>> + Send + 'static,
{
let database = crate::database_capability::request_storage_layouts()
.map_err(|_| saddle_admission::AdmissionError::SizeOverflow)?;
let context = consumer_storage::shared_layout::<ReservedTaskContext>()
.map_err(|_| saddle_admission::AdmissionError::SizeOverflow)?;
let call_counter = consumer_storage::shared_layout::<std::sync::atomic::AtomicU64>()
.map_err(|_| saddle_admission::AdmissionError::SizeOverflow)?;
let indirect = [
(database.state.allocation, 1),
(database.cancellation_notify.allocation, 1),
(database.outbound_retentions.allocation, 1),
(context.allocation, 1),
(
std::alloc::Layout::new::<tokio::sync::futures::OwnedNotified>(),
1,
),
(call_counter.allocation, 1),
];
let body = body_layout::<B, C, Dispatch, DispatchFuture>();
let bytes = saddle_runtime::request_task::reserved::dispatch_storage_bytes_for::<
NormalInput<B, C, Dispatch>,
(),
_,
>(&factory::<B, C, Dispatch, DispatchFuture>, body, &indirect)?;
fn output<A, R>(_: impl FnOnce(A) -> R) -> usize {
std::mem::size_of::<R>()
}
let ingress = output(
|(factory, dispatch, admission, task): (
&'static mut consumer_storage::NormalFactory<B, C, Dispatch>,
&'static mut Option<saddle_runtime::profusegw::ProfuseGwManagedDispatch>,
&'static mut Option<saddle_runtime::profusegw::ProfuseGwAdmissionEvent>,
&'static mut ReservedTaskContext,
)| factory.reserved_ingress(dispatch, admission, task),
);
Ok((
std::alloc::Layout::new::<Owner<B, C, Dispatch>>(),
body,
bytes,
[
ingress,
std::mem::size_of::<DispatchFuture>(),
std::mem::size_of::<crate::database_capability::DatabaseRequest>(),
std::mem::size_of::<saddle_observability::RequestDiagnosticScope<'static>>(),
database.state.allocation.size(),
],
))
}