use super::*;
use saddle_admission::{AdmissionError, InputDocument, InputKind, InputValue};
#[derive(Debug)]
pub enum ManagedIngressError { Codec(CodecError), Storage(AdmissionError) }
impl ManagedIngressError {
pub fn http_status(&self) -> u16 { match self { Self::Codec(e) => e.http_status, Self::Storage(_) => 503 } }
}
impl std::fmt::Display for ManagedIngressError {
fn fmt(&self,f:&mut std::fmt::Formatter<'_>)->std::fmt::Result { write!(f,"{self:?}") }
}
impl std::error::Error for ManagedIngressError {
fn source(&self)->Option<&(dyn std::error::Error+'static)> { match self { Self::Storage(e)=>Some(e), _=>None } }
}
impl From<CodecError> for ManagedIngressError { fn from(e:CodecError)->Self { Self::Codec(e) } }
impl From<AdmissionError> for ManagedIngressError { fn from(e:AdmissionError)->Self { Self::Storage(e) } }
fn shape(value:InputValue<'_>, keys:&[&str], optional:Option<&str>)->Result<(),CodecError> {
let mut seen=0u64;
for (key,_) in value.fields().ok_or(error(400,"INVALID_JSON_ENVELOPE"))? {
let i=keys.iter().position(|k|*k==key).ok_or(error(400,"INVALID_JSON_ENVELOPE"))?;
if seen&(1<<i)!=0 { return Err(error(400,"INVALID_JSON_ENVELOPE")); }
seen|=1<<i;
}
for (i,key) in keys.iter().enumerate() { if seen&(1<<i)==0 && Some(*key)!=optional { return Err(error(400,"INVALID_JSON_ENVELOPE")); } }
Ok(())
}
fn string<'a>(v:InputValue<'a>,key:&str)->Result<&'a str,CodecError> {
let value=v.field(key).ok_or(error(400,"INVALID_JSON_ENVELOPE"))?;
if value.kind()!=InputKind::String { return Err(error(400,"INVALID_JSON_ENVELOPE")); }
Ok(value.text().unwrap())
}
fn unique_json(value: InputValue<'_>) -> Result<(), CodecError> {
if let Some(fields) = value.fields() {
for (index, (key, child)) in fields.enumerate() {
if value.fields().unwrap().take(index).any(|(previous, _)| previous == key) {
return Err(error(400, "INVALID_JSON_ENVELOPE"));
}
unique_json(child)?;
}
} else if let Some(elements) = value.elements() {
for child in elements { unique_json(child)?; }
}
Ok(())
}
fn extensible_object(value: InputValue<'_>) -> Result<(), CodecError> {
if value.kind() != InputKind::Object { return Err(error(400, "INVALID_JSON_ENVELOPE")); }
unique_json(value)
}
impl ProfuseGwListenerAdapter {
#[doc(hidden)]
pub fn validate_managed_transport(&self,method:&str,path:&str,content_type:&str,identity:&IngressIdentity,body:&[u8])->Result<(),ManagedIngressError> {
validate_identity(identity)?; validate_transport(method,path,content_type,body)?; Ok(())
}
#[doc(hidden)]
pub fn accept_managed(&self,method:&str,path:&str,content_type:&str,identity:&IngressIdentity,body:&[u8],document:InputDocument)->Result<AcceptedIngress,ManagedIngressError> {
validate_identity(identity)?;
validate_transport(method,path,content_type,body)?;
let root=document.root();
shape(root,&["target","profuseGwContext","requestData","profuseContext"],Some("profuseContext"))?;
let target=root.field("target").unwrap();
shape(target,&["app","interfaceId"],None)?;
let app=string(target,"app")?;
let interface=string(target,"interfaceId")?;
let context=root.field("profuseGwContext").unwrap();
shape(context,&["userInfo","traceInfo","ldcInfo"],None)?;
let user=context.field("userInfo").unwrap(); shape(user,&["userId"],None)?;
let user=string(user,"userId")?;
let trace=context.field("traceInfo").unwrap();
let ldc=context.field("ldcInfo").unwrap();
let profuse=root.field("profuseContext");
if trace.source_range().len() + ldc.source_range().len() + profuse.map_or(0, |v| v.source_range().len()) > crate::grouped_context::LIMIT {
return Err(error(413,"CONTEXT_TOO_LARGE").into());
}
extensible_object(trace)?;
let trace_id=if trace.field("traceId").is_some() { string(trace,"traceId")? } else { "" };
let rpc=string(trace,"rpcId")?;
extensible_object(ldc)?;
let zone=string(ldc,"zone")?; let idc=string(ldc,"idc")?; let env=string(ldc,"env")?;
for value in [app,interface,user,rpc,zone,idc,env] { validate_id(value)?; }
if !trace_id.is_empty() { validate_trace_id(trace_id)?; }
if root.field("requestData").unwrap().kind()!=InputKind::Object { return Err(error(400,"INVALID_REQUEST_DATA").into()); }
if app!=self.application.as_str() { return Err(error(404,"APPLICATION_NOT_FOUND").into()); }
let profuse = root.field("profuseContext");
if let Some(value) = profuse { extensible_object(value)?; }
let raw = |value: InputValue<'_>| -> Result<&str, CodecError> {
let bytes = body.get(value.source_range()).ok_or(error(400,"INVALID_JSON_ENVELOPE"))?;
std::str::from_utf8(bytes).map_err(|_| error(400,"INVALID_JSON_ENVELOPE"))
};
let trace_json = raw(trace)?;
let ldc_json = raw(ldc)?;
let profuse_json = profuse.map(raw).transpose()?;
if trace_json.len() + ldc_json.len() + profuse_json.map_or(0, str::len) > 4096 {
return Err(error(413,"CONTEXT_TOO_LARGE").into());
}
let context_json = root.framework_shared_input(|b| Ok(IngressContextJson {
trace_info: b.copy_text(&[trace_json])?,
rpc_value: {
let range = trace.field("rpcId").unwrap().source_range();
range.start - trace.source_range().start..range.end - trace.source_range().start
},
trace_value: trace.field("traceId").map(|value| {
let range = value.source_range();
range.start - trace.source_range().start..range.end - trace.source_range().start
}),
ldc_info: b.copy_text(&[ldc_json])?,
profuse_context: profuse_json.map(|value| b.copy_text(&[value])).transpose()?,
}))?;
let metadata=root.framework_input(|b| Ok(IngressMetadata {
identity: IngressIdentity { request_id:b.copy_text(&[&identity.request_id])?, call_id:b.copy_text(&[&identity.call_id])?, deadline_unix_ms:identity.deadline_unix_ms },
interface_id:b.copy_text(&[interface])?, user_id:b.copy_text(&[user])?,
trace_id:if trace_id.is_empty() { b.copy_text(&["saddle-",&identity.request_id])? } else { b.copy_text(&[trace_id])? },
rpc_id:b.copy_text(&[rpc])?, zone:b.copy_text(&[zone])?, idc:b.copy_text(&[idc])?, env:b.copy_text(&[env])?,
}))?;
Ok(AcceptedIngress { metadata:AcceptedMetadata::Managed(metadata),context_json:Some(IngressContextOwner(ContextOwner::Managed(context_json))),request_data:AcceptedRequestData::Managed(document) })
}
}