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();
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))?;
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,
})
}
}