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),
)
}
}