1use std::pin::Pin;
2use std::sync::Arc;
3
4use bytes::Bytes;
5use futures_util::Stream;
6use unb_core::{Envelope, ErrorCode};
7
8use crate::layer::{Next, Origin, ServiceBody};
9use crate::node::{Node, NodeSnapshot};
10use crate::service::Operation;
11
12pub type EventStream = Pin<Box<dyn Stream<Item = Result<Bytes, HandlerError>> + Send>>;
13
14#[derive(Debug, thiserror::Error, schemars::JsonSchema)]
15#[error("{code:?}: {message}")]
16pub struct HandlerError {
17 pub code: ErrorCode,
18 pub message: String,
19}
20
21impl HandlerError {
22 pub fn new(code: ErrorCode, message: impl Into<String>) -> HandlerError {
23 HandlerError {
24 code,
25 message: message.into(),
26 }
27 }
28}
29
30impl Node {
31 pub(crate) async fn run_service(
32 self: &Arc<Self>,
33 snapshot: Arc<NodeSnapshot>,
34 mut request: http::Request<Bytes>,
35 origin: Origin,
36 ) -> Option<Result<http::Response<ServiceBody>, HandlerError>> {
37 let subject = Envelope::subject_of(request.uri());
38 let Some(operation) = Operation::of(&request) else {
39 return Some(Err(HandlerError::new(
40 ErrorCode::Protocol,
41 format!("method {:?} maps to no operation", request.method()),
42 )));
43 };
44 let entry = snapshot.services.get(&subject).cloned()?;
45 let compiled = match operation {
46 Operation::Unary => entry.unary.clone(),
47 Operation::Streaming => entry.streaming.clone(),
48 };
49 let Some(compiled) = compiled else {
50 return Some(Err(HandlerError::new(
51 ErrorCode::Protocol,
52 format!("subject {subject:?} does not serve {operation:?} operations"),
53 )));
54 };
55 request.extensions_mut().insert(origin);
56 request.extensions_mut().insert(Arc::downgrade(self));
57 request.extensions_mut().insert(snapshot);
58 Some(
59 Next::root(compiled.layers.clone(), compiled.call.clone())
60 .run(request)
61 .await,
62 )
63 }
64
65 pub(crate) fn inbound_request(
66 envelope: &Envelope,
67 ) -> Result<http::Request<Bytes>, HandlerError> {
68 envelope
69 .to_local_request()
70 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))
71 }
72
73 pub(crate) fn local_request(
74 &self,
75 kind: unb_core::Kind,
76 subject: &str,
77 payload: Bytes,
78 headers: serde_json::Map<String, serde_json::Value>,
79 ) -> Result<http::Request<Bytes>, HandlerError> {
80 let envelope = Envelope {
81 v: unb_core::PROTOCOL_VERSION,
82 id: String::new(),
83 target: self.identity().node_id.clone(),
84 subject: subject.to_string(),
85 kind,
86 corr: None,
87 seq: None,
88 hops: None,
89 body_token: None,
90 payload: Bytes::new(),
91 path: Vec::new(),
92 headers,
93 };
94 let request = envelope
95 .to_local_request()
96 .map_err(|error| HandlerError::new(ErrorCode::InvalidInput, error.to_string()))?;
97 Ok(request.map(|_| payload))
98 }
99
100 pub(crate) fn teach_unknown_subject(snapshot: &NodeSnapshot, subject: &str) -> HandlerError {
101 let known_refs: Vec<&str> = snapshot.services.keys().map(String::as_str).collect();
102 HandlerError::new(
103 ErrorCode::UnknownSubject,
104 unb_core::teach_unknown("subject", subject, &known_refs),
105 )
106 }
107
108 pub(crate) fn teach_unknown_target(snapshot: &NodeSnapshot, target: &str) -> HandlerError {
109 let known = snapshot.node_core.reachable_names();
110 let known_refs: Vec<&str> = known.iter().map(String::as_str).collect();
111 HandlerError::new(
112 ErrorCode::PeerUnreachable,
113 unb_core::teach_unknown("node", target, &known_refs),
114 )
115 }
116}