saddle-boundary 0.3.34

Saddle 0.3 ProfuseContract unary boundary transport
//! Fixed envelope validation and prepaid metadata; no serde scratch or Value.
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())
}
// Scan sibling keys without allocating a second object tree. The input parser
// has already bounded nesting and prepaid all decoded key storage.
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) })
    }
}