saddle-framework 0.3.25

The single business-facing facade for Saddle applications
//! Same production parser/adapter and NormalFactory captures, with the AR
//! event remaining in the external owner until identity is fully validated.
use super::*;
use saddle_core::request_context::{ContextFact, ContextLabel, RequestIdentityGroup};
use saddle_runtime::request_task::reserved::{ReservedRequestFailure, ReservedTaskContext};

pub(super) struct PreparedIngress {
    pub accepted: saddle_boundary::ingress::AcceptedIngress,
    pub call: saddle_observability::ActiveCall,
    pub event: saddle_observability::EventContext,
    pub admission_submission: saddle_observability::AdmissionCapacitySubmission,
}
pub(super) struct IngressFailure {
    pub status: u16,
    pub source: ReservedRequestFailure,
}
#[derive(Debug)]
struct Description<'a, T: std::fmt::Debug>(&'a T);
impl<T: std::fmt::Debug> std::fmt::Display for Description<'_, T> {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        write!(f, "{:?}", self.0)
    }
}
impl<B, C, Dispatch> consumer_storage::NormalFactory<B, C, Dispatch> {
    pub(super) async fn reserved_ingress(
        &mut self,
        dispatch: &mut Option<saddle_runtime::profusegw::ProfuseGwManagedDispatch>,
        admission: &mut Option<saddle_runtime::profusegw::ProfuseGwAdmissionEvent>,
        task: &mut ReservedTaskContext,
    ) -> std::result::Result<PreparedIngress, IngressFailure> {
        use saddle_boundary::source_diagnostics::IoOperation;
        use saddle_core::request_context::RequestViewPhase;
        let output = self.diagnostic_handle.as_ref();
        let io = |task: &ReservedTaskContext, error: std::io::Error, operation| {
            let (_, source) = saddle_boundary::reserved_diagnostics::capture_io(
                &task.view(),
                output,
                operation,
                &error,
            )
            .into_parts();
            IngressFailure {
                status: if error.kind() == std::io::ErrorKind::PermissionDenied {
                    401
                } else {
                    400
                },
                source,
            }
        };
        let fault = |task: &ReservedTaskContext, value: &dyn std::fmt::Debug, status| {
            use saddle_core::{
                BoundedDiagnostic, BoundedDiagnosticCause, CaptureSite, DiagnosticCategory,
                DiagnosticCode, DiagnosticStage,
            };
            let code = DiagnosticCode::new("service.ingress.validation").unwrap();
            let source = task.view().source_description(
                &Description(&value),
                BoundedDiagnostic::capture(
                    DiagnosticCategory::UnexpectedError,
                    CaptureSite::FirstObserved,
                    BoundedDiagnosticCause::new(DiagnosticStage::RequestDecode, code),
                ),
                code,
                output,
                saddle_observability::root_diagnostic::RootRequestEvent::Ingress,
                Default::default(),
            );
            IngressFailure { status, source }
        };
        let view = task
            .view()
            .with_phase(RequestViewPhase::Reading)
            .map_err(|e| fault(task, &e, 500))?;
        task.observe_view(view)
            .unwrap_or_else(|_| unreachable!("same original request"));
        let initial_deadline = dispatch
            .as_ref()
            .expect("original dispatch")
            .deadline_unix_ms();
        // Original admission is untouched while either real socket read is Pending.
        let head = tokio::select! {
            biased;
            ()=dispatch.as_mut().unwrap().deadline_elapsed()=>return Err(io(task,std::io::ErrorKind::TimedOut.into(),IoOperation::ReadHead)),
            value=read_request_head(&mut self.socket,initial_deadline,self.ingress_token.as_deref())=>value.map_err(|e|io(task,e,IoOperation::ReadHead))?,
        };
        dispatch
            .as_mut()
            .unwrap()
            .bind_ingress_deadline(head.deadline_unix_ms)
            .map_err(|e| fault(task, &e, 400))?;
        let request = tokio::select! {
            biased;
            ()=dispatch.as_mut().unwrap().deadline_elapsed()=>return Err(io(task,std::io::ErrorKind::TimedOut.into(),IoOperation::ReadBody)),
            value=read_request_body(&mut self.socket,head)=>value.map_err(|e|io(task,e,IoOperation::ReadBody))?,
        };
        let mut parser_source = None;
        let accepted = self
            .adapter
            .accept_observed(
                &request.method,
                &request.path,
                &request.content_type,
                request.identity,
                &request.body,
                |original| {
                    parser_source = Some(task.view().source_error(
                        original,
                        output,
                        saddle_core::DiagnosticStage::RequestDecode,
                        saddle_observability::root_diagnostic::RootRequestEvent::Ingress,
                        Default::default(),
                    ));
                },
            )
            .map_err(|e| match parser_source.take() {
                Some(source) => IngressFailure {
                    status: e.http_status,
                    source,
                },
                None => fault(task, &e, e.http_status),
            })?;
        let rpc = saddle_core::RpcCorrelationId::new(&accepted.rpc_id)
            .ok_or_else(|| fault(task, &"validated RPC unavailable", 400))?;
        let (call, _) = self
            .observer
            .start_external_call_with_rpc(
                self.adapter.application(),
                "profusegw",
                "profusegw",
                accepted.interface_id.clone(),
                Some(&accepted.trace_id),
                rpc,
            )
            .map_err(|e| fault(task, &e, 400))?;
        let event = saddle_observability::EventContext::new(
            saddle_observability::RequestIdentity::new(accepted.identity.request_id.clone())
                .map_err(|e| fault(task, &e, 400))?,
            saddle_observability::RouteIdentity::new(accepted.interface_id.clone())
                .map_err(|e| fault(task, &e, 400))?,
            1,
        )
        .map_err(|e| fault(task, &e, 400))?;
        let zone = ContextFact::Present(
            ContextLabel::checked(&accepted.zone).map_err(|e| fault(task, &e, 400))?,
        );
        let identity = RequestIdentityGroup::from_validated(
            call.context(),
            &accepted.identity.request_id,
            &accepted.interface_id,
            1,
            zone,
        )
        .map_err(|e| fault(task, &e, 400))?;
        task.publish(identity).map_err(|e| fault(task, &e, 500))?;
        // No await between taking the actual event and consuming it exactly once.
        let event_receipt = admission
            .take()
            .expect("AR preserves unique original admission");
        let (bound, admission_submission) = dispatch.take().unwrap().bind_reserved_observation(
            self.observer.clone(),
            call.context().clone(),
            event.clone(),
            event_receipt,
            &task.view(),
            output,
        );
        *dispatch = Some(bound);
        Ok(PreparedIngress {
            accepted,
            call,
            event,
            admission_submission,
        })
    }
}