unb-server 2.0.3

unb inbound server: Node, request/subscribe handlers, catalog, relay orchestration, accept
Documentation
use std::pin::Pin;
use std::sync::Arc;

use bytes::Bytes;
use futures_util::Stream;
use unb_core::{Envelope, ErrorCode};

use crate::layer::{Next, Origin, ServiceBody};
use crate::node::{Node, NodeSnapshot};
use crate::service::Operation;

pub type EventStream = Pin<Box<dyn Stream<Item = Result<Bytes, HandlerError>> + Send>>;

#[derive(Debug, thiserror::Error, schemars::JsonSchema)]
#[error("{code:?}: {message}")]
pub struct HandlerError {
    pub code: ErrorCode,
    pub message: String,
}

impl HandlerError {
    pub fn new(code: ErrorCode, message: impl Into<String>) -> HandlerError {
        HandlerError {
            code,
            message: message.into(),
        }
    }
}

impl Node {
    pub(crate) async fn run_service(
        self: &Arc<Self>,
        snapshot: Arc<NodeSnapshot>,
        mut request: http::Request<Bytes>,
        origin: Origin,
    ) -> Option<Result<http::Response<ServiceBody>, HandlerError>> {
        let subject = Envelope::subject_of(request.uri());
        let Some(operation) = Operation::of(&request) else {
            return Some(Err(HandlerError::new(
                ErrorCode::Protocol,
                format!("method {:?} maps to no operation", request.method()),
            )));
        };
        let entry = snapshot.services.get(&subject).cloned()?;
        let compiled = match operation {
            Operation::Unary => entry.unary.clone(),
            Operation::Streaming => entry.streaming.clone(),
        };
        let Some(compiled) = compiled else {
            return Some(Err(HandlerError::new(
                ErrorCode::Protocol,
                format!("subject {subject:?} does not serve {operation:?} operations"),
            )));
        };
        request.extensions_mut().insert(origin);
        request.extensions_mut().insert(Arc::downgrade(self));
        request.extensions_mut().insert(snapshot);
        Some(
            Next::root(compiled.layers.clone(), compiled.call.clone())
                .run(request)
                .await,
        )
    }

    pub(crate) fn inbound_request(
        envelope: &Envelope,
    ) -> Result<http::Request<Bytes>, HandlerError> {
        envelope
            .to_local_request()
            .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))
    }

    pub(crate) fn local_request(
        &self,
        kind: unb_core::Kind,
        subject: &str,
        payload: Bytes,
        headers: serde_json::Map<String, serde_json::Value>,
    ) -> Result<http::Request<Bytes>, HandlerError> {
        let envelope = Envelope {
            v: unb_core::PROTOCOL_VERSION,
            id: String::new(),
            target: self.identity().node_id.clone(),
            subject: subject.to_string(),
            kind,
            corr: None,
            seq: None,
            hops: None,
            body_token: None,
            payload: Bytes::new(),
            path: Vec::new(),
            headers,
        };
        let request = envelope
            .to_local_request()
            .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
        Ok(request.map(|_| payload))
    }

    pub(crate) fn teach_unknown_subject(snapshot: &NodeSnapshot, subject: &str) -> HandlerError {
        let known_refs: Vec<&str> = snapshot.services.keys().map(String::as_str).collect();
        HandlerError::new(
            ErrorCode::UnknownSubject,
            unb_core::teach_unknown("subject", subject, &known_refs),
        )
    }

    pub(crate) fn teach_unknown_target(snapshot: &NodeSnapshot, target: &str) -> HandlerError {
        let known = snapshot.node_core.reachable_names();
        let known_refs: Vec<&str> = known.iter().map(String::as_str).collect();
        HandlerError::new(
            ErrorCode::PeerUnreachable,
            unb_core::teach_unknown("node", target, &known_refs),
        )
    }
}