1use bytes::{BufMut, Bytes, BytesMut};
2use serde::{Deserialize, Serialize};
3use serde_json::Value;
4
5use crate::{CoreError, TargetPath};
6
7pub const PROTOCOL_VERSION: u16 = 1;
8pub const DEFAULT_HOPS: u8 = 8;
9
10const CONTROL_TARGET: &str = "*";
11const MAX_HEADERS: usize = 64;
12const MAX_HEAD_BYTES: usize = 64 * 1024;
13
14pub const UNB_VERSION: &str = "unb-version";
15pub const UNB_KIND: &str = "unb-kind";
16pub const UNB_ID: &str = "unb-id";
17pub const UNB_CORR: &str = "unb-corr";
18pub const UNB_SEQ: &str = "unb-seq";
19pub const UNB_HOPS: &str = "unb-hops";
20pub const UNB_PATH: &str = "unb-path";
21pub const UNB_CODE: &str = "unb-code";
22pub const UNB_BODY: &str = "unb-body";
23
24#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
25#[serde(rename_all = "snake_case")]
26pub enum Kind {
27 Hello,
28 Welcome,
29 Request,
30 Response,
31 Subscribe,
32 Event,
33 Channel,
34 Cancel,
35 Error,
36 Ping,
37 Pong,
38 Discover,
39 Identify,
40 IdentityAccepted,
41 RouteSnapshot,
42 RouteDelta,
43 RouteAck,
44}
45
46impl Kind {
47 pub fn token(self) -> &'static str {
48 match self {
49 Kind::Hello => "hello",
50 Kind::Welcome => "welcome",
51 Kind::Request => "request",
52 Kind::Response => "response",
53 Kind::Subscribe => "subscribe",
54 Kind::Event => "event",
55 Kind::Channel => "channel",
56 Kind::Cancel => "cancel",
57 Kind::Error => "error",
58 Kind::Ping => "ping",
59 Kind::Pong => "pong",
60 Kind::Discover => "discover",
61 Kind::Identify => "identify",
62 Kind::IdentityAccepted => "identity_accepted",
63 Kind::RouteSnapshot => "route_snapshot",
64 Kind::RouteDelta => "route_delta",
65 Kind::RouteAck => "route_ack",
66 }
67 }
68
69 pub fn from_token(token: &str) -> Option<Kind> {
70 Some(match token {
71 "hello" => Kind::Hello,
72 "welcome" => Kind::Welcome,
73 "request" => Kind::Request,
74 "response" => Kind::Response,
75 "subscribe" => Kind::Subscribe,
76 "event" => Kind::Event,
77 "channel" => Kind::Channel,
78 "cancel" => Kind::Cancel,
79 "error" => Kind::Error,
80 "ping" => Kind::Ping,
81 "pong" => Kind::Pong,
82 "discover" => Kind::Discover,
83 "identify" => Kind::Identify,
84 "identity_accepted" => Kind::IdentityAccepted,
85 "route_snapshot" => Kind::RouteSnapshot,
86 "route_delta" => Kind::RouteDelta,
87 "route_ack" => Kind::RouteAck,
88 _ => return None,
89 })
90 }
91
92 pub fn is_application_request(self) -> bool {
93 matches!(
94 self,
95 Kind::Request | Kind::Subscribe | Kind::Channel | Kind::Discover
96 )
97 }
98
99 pub fn is_application_response(self) -> bool {
100 matches!(self, Kind::Response | Kind::Event | Kind::Error)
101 }
102}
103
104#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
105pub struct Envelope {
106 pub v: u16,
107 pub id: String,
108 #[serde(default, skip_serializing_if = "String::is_empty")]
109 pub target: String,
110 #[serde(default, skip_serializing_if = "String::is_empty")]
111 pub subject: String,
112 pub kind: Kind,
113 #[serde(default, skip_serializing_if = "Option::is_none")]
114 pub corr: Option<String>,
115 #[serde(default, skip_serializing_if = "Option::is_none")]
116 pub seq: Option<u64>,
117 #[serde(default, skip_serializing_if = "Option::is_none")]
118 pub hops: Option<u8>,
119 #[serde(default, skip_serializing_if = "Option::is_none")]
120 pub body_token: Option<String>,
121 #[serde(skip)]
122 pub payload: Bytes,
123 #[serde(default, skip_serializing_if = "Vec::is_empty")]
124 pub path: Vec<String>,
125 #[serde(flatten)]
126 pub headers: serde_json::Map<String, Value>,
127}
128
129impl Envelope {
130 pub fn decode(frame: Bytes) -> Result<Envelope, CoreError> {
131 let (mut envelope, head_len) = Envelope::decode_head(&frame)?;
132 envelope.payload = frame.slice(head_len..);
133 Ok(envelope)
134 }
135
136 pub fn decode_head(frame: &[u8]) -> Result<(Envelope, usize), CoreError> {
137 if frame.len() > MAX_HEAD_BYTES && !Envelope::head_within_bounds(frame) {
138 return Err(CoreError::Malformed(
139 "header section exceeds the maximum head size".into(),
140 ));
141 }
142 let mut envelope = Envelope {
143 v: PROTOCOL_VERSION,
144 id: String::new(),
145 target: String::new(),
146 subject: String::new(),
147 kind: Kind::Request,
148 corr: None,
149 seq: None,
150 hops: None,
151 body_token: None,
152 payload: Bytes::new(),
153 path: Vec::new(),
154 headers: serde_json::Map::new(),
155 };
156 let mut storage = [httparse::EMPTY_HEADER; MAX_HEADERS];
157 let is_status = frame.starts_with(b"HTTP/");
158 let (head_len, status, target) = if is_status {
159 let mut response = httparse::Response::new(&mut storage);
160 let head_len = match response.parse(frame).map_err(Envelope::parse_error)? {
161 httparse::Status::Complete(head_len) => head_len,
162 httparse::Status::Partial => {
163 return Err(CoreError::Malformed("frame head is incomplete".into()))
164 }
165 };
166 let code = response
167 .code
168 .ok_or_else(|| CoreError::Malformed("status line carries no status".into()))?;
169 (head_len, Some(code), None)
170 } else {
171 let mut request = httparse::Request::new(&mut storage);
172 let head_len = match request.parse(frame).map_err(Envelope::parse_error)? {
173 httparse::Status::Complete(head_len) => head_len,
174 httparse::Status::Partial => {
175 return Err(CoreError::Malformed("frame head is incomplete".into()))
176 }
177 };
178 if request.method != Some("POST") {
179 return Err(CoreError::Malformed(
180 "application and control frames use POST".into(),
181 ));
182 }
183 let target = request
184 .path
185 .ok_or_else(|| CoreError::Malformed("request line carries no target".into()))?;
186 (head_len, None, Some(target.to_string()))
187 };
188 let az_kind = envelope.absorb_wire_headers(&storage, is_status)?;
189 if let Some(code) = status {
190 if az_kind.is_some() {
191 return Err(CoreError::Malformed(
192 "status-line frames carry no unb-kind".into(),
193 ));
194 }
195 envelope.kind = match code {
196 200 if envelope.seq.is_some() => Kind::Event,
197 200 => Kind::Response,
198 400..=599 => Kind::Error,
199 other => {
200 return Err(CoreError::Malformed(format!(
201 "status {other} maps to no frame kind"
202 )));
203 }
204 };
205 } else {
206 let target = target.expect("request path present");
207 if target.contains(['?', '#']) {
208 return Err(CoreError::Malformed(
209 "targets must not contain uri query or fragment delimiters".into(),
210 ));
211 }
212 if target == CONTROL_TARGET {
213 let control = az_kind
214 .ok_or_else(|| CoreError::Malformed("control frames carry unb-kind".into()))?;
215 if control.is_application_request() || control.is_application_response() {
216 return Err(CoreError::Malformed(format!(
217 "{control:?} is not a control kind"
218 )));
219 }
220 envelope.kind = control;
221 } else if target.starts_with('/') {
222 let kind = az_kind.unwrap_or(Kind::Request);
223 if !kind.is_application_request() {
224 return Err(CoreError::Malformed(format!(
225 "{kind:?} is not an application request kind"
226 )));
227 }
228 let target_path = if kind == Kind::Discover {
229 TargetPath::parse_discovery(&target)?
230 } else {
231 TargetPath::parse_application(&target)?
232 };
233 envelope.target = target_path.target().to_owned();
234 envelope.subject = target_path.subject().to_owned();
235 envelope.kind = kind;
236 } else {
237 return Err(CoreError::Malformed(
238 "request targets are origin-form and start with a slash".into(),
239 ));
240 }
241 }
242 Ok((envelope, head_len))
243 }
244
245 pub fn encode(&self) -> Bytes {
246 use std::fmt::Write;
247 let mut head = String::with_capacity(64 + self.headers.len() * 32);
248 if self.kind.is_application_response() {
249 let (status, reason) = match self.kind {
250 Kind::Error => match self.error_code() {
251 Some(code) => (code.status().as_u16(), code.token()),
252 None => (500, "INTERNAL"),
253 },
254 _ => (200, "OK"),
255 };
256 write!(head, "HTTP/1.1 {status} {reason}\r\n")
257 .expect("writing to a string cannot fail");
258 } else if self.kind.is_application_request() {
259 let target_path = if self.kind == Kind::Discover {
260 TargetPath::discovery(self.target.clone())
261 } else {
262 TargetPath::application(self.target.clone(), self.subject.clone())
263 }
264 .expect("application envelopes carry a valid target path");
265 write!(head, "POST {target_path} HTTP/1.1\r\n")
266 .expect("writing to a string cannot fail");
267 } else {
268 write!(head, "POST {CONTROL_TARGET} HTTP/1.1\r\n")
269 .expect("writing to a string cannot fail");
270 }
271 let line = |head: &mut String, name: &str, value: &str| {
272 debug_assert!(
273 !value.contains(['\r', '\n']),
274 "header values are single-line strings"
275 );
276 head.push_str(name);
277 head.push_str(": ");
278 head.push_str(value);
279 head.push_str("\r\n");
280 };
281 line(&mut head, UNB_VERSION, &self.v.to_string());
282 if !self.id.is_empty() {
283 line(&mut head, UNB_ID, &self.id);
284 }
285 if let Some(corr) = &self.corr {
286 line(&mut head, UNB_CORR, corr);
287 }
288 if let Some(seq) = self.seq {
289 line(&mut head, UNB_SEQ, &seq.to_string());
290 }
291 if let Some(hops) = self.hops {
292 line(&mut head, UNB_HOPS, &hops.to_string());
293 }
294 if let Some(token) = &self.body_token {
295 line(&mut head, UNB_BODY, token);
296 }
297 if !self.path.is_empty() {
298 line(&mut head, UNB_PATH, &self.path.join(","));
299 }
300 if !self.kind.is_application_response() && self.kind != Kind::Request {
301 line(&mut head, UNB_KIND, self.kind.token());
302 }
303 if self.kind == Kind::Error {
304 if let Some(code) = self.payload_json()["code"].as_str() {
305 line(&mut head, UNB_CODE, code);
306 }
307 }
308 for (name, value) in &self.headers {
309 let value = value
310 .as_str()
311 .expect("application headers carry string values on the wire");
312 line(&mut head, name, value);
313 }
314 head.push_str("\r\n");
315 let mut frame = BytesMut::with_capacity(head.len() + self.payload.len());
316 frame.put_slice(head.as_bytes());
317 frame.put_slice(&self.payload);
318 frame.freeze()
319 }
320
321 fn head_within_bounds(frame: &[u8]) -> bool {
322 frame
323 .windows(4)
324 .take(MAX_HEAD_BYTES)
325 .position(|window| window == b"\r\n\r\n")
326 .is_some_and(|position| position + 4 <= MAX_HEAD_BYTES)
327 }
328
329 fn parse_error(error: httparse::Error) -> CoreError {
330 CoreError::Malformed(error.to_string())
331 }
332
333 fn absorb_wire_headers(
334 &mut self,
335 headers: &[httparse::Header],
336 is_status: bool,
337 ) -> Result<Option<Kind>, CoreError> {
338 let mut az_kind = None;
339 let mut az_version_seen = false;
340 for header in headers {
341 if header.name.is_empty() {
342 continue;
343 }
344 let name = header.name.to_ascii_lowercase();
345 let value = std::str::from_utf8(header.value)
346 .map_err(|source| CoreError::Malformed(source.to_string()))?
347 .trim_start_matches(' ');
348 let duplicate = || CoreError::Malformed(format!("{name}: duplicate header"));
349 match name.as_str() {
350 "content-length" | "transfer-encoding" => {
351 return Err(CoreError::Malformed(format!(
352 "{name}: message-carrier frames carry no explicit body framing"
353 )));
354 }
355 UNB_VERSION => {
356 if az_version_seen {
357 return Err(duplicate());
358 }
359 az_version_seen = true;
360 self.v = value
361 .parse()
362 .map_err(|_| CoreError::Malformed(format!("{UNB_VERSION}: {value:?}")))?;
363 }
364 UNB_ID => {
365 if !self.id.is_empty() {
366 return Err(duplicate());
367 }
368 self.id = value.to_string();
369 }
370 UNB_CORR => {
371 if self.corr.is_some() {
372 return Err(duplicate());
373 }
374 self.corr = Some(value.to_string());
375 }
376 UNB_SEQ => {
377 if self.seq.is_some() {
378 return Err(duplicate());
379 }
380 self.seq = Some(
381 value
382 .parse()
383 .map_err(|_| CoreError::Malformed(format!("{UNB_SEQ}: {value:?}")))?,
384 );
385 }
386 UNB_HOPS => {
387 if self.hops.is_some() {
388 return Err(duplicate());
389 }
390 self.hops = Some(
391 value
392 .parse()
393 .map_err(|_| CoreError::Malformed(format!("{UNB_HOPS}: {value:?}")))?,
394 );
395 }
396 UNB_PATH => {
397 if !self.path.is_empty() {
398 return Err(duplicate());
399 }
400 self.path = value.split(',').map(str::to_string).collect();
401 }
402 UNB_BODY => {
403 if self.body_token.is_some() {
404 return Err(duplicate());
405 }
406 if value.is_empty() || value.len() > 256 {
407 return Err(CoreError::Malformed(format!(
408 "{UNB_BODY}: token length out of bounds"
409 )));
410 }
411 self.body_token = Some(value.to_string());
412 }
413 UNB_KIND => {
414 if az_kind.is_some() {
415 return Err(duplicate());
416 }
417 az_kind = Some(Kind::from_token(value).ok_or_else(|| {
418 CoreError::Malformed(format!("{UNB_KIND}: {value:?} is not a frame kind"))
419 })?);
420 }
421 UNB_CODE if is_status => {}
422 other if other.starts_with("unb-") => {
423 return Err(CoreError::Malformed(format!(
424 "{other}: unknown reserved header"
425 )));
426 }
427 _ => {
428 if self
429 .headers
430 .insert(name.clone(), Value::String(value.to_string()))
431 .is_some()
432 {
433 return Err(duplicate());
434 }
435 }
436 }
437 }
438 Ok(az_kind)
439 }
440
441 fn error_code(&self) -> Option<crate::ErrorCode> {
442 serde_json::from_value(self.payload_json()["code"].clone()).ok()
443 }
444
445 #[inline]
446 pub fn payload_json(&self) -> Value {
447 if self.payload.is_empty() {
448 return Value::Null;
449 }
450 serde_json::from_slice(&self.payload).unwrap_or(Value::Null)
451 }
452
453 #[inline]
454 pub fn encode_payload(value: &Value) -> Bytes {
455 if value.is_null() {
456 return Bytes::new();
457 }
458 serde_json::to_vec(value)
459 .expect("payload is plain json data")
460 .into()
461 }
462
463 #[inline]
464 pub fn parse_payload<T: serde::de::DeserializeOwned>(&self) -> Result<T, CoreError> {
465 serde_json::from_slice(&self.payload)
466 .map_err(|source| CoreError::Malformed(source.to_string()))
467 }
468
469 pub fn to_request(&self) -> Result<http::Request<Bytes>, CoreError> {
470 if !self.kind.is_application_request() {
471 return Err(CoreError::Malformed(format!(
472 "{:?} is not an application request kind",
473 self.kind
474 )));
475 }
476 let target_path = if self.kind == Kind::Discover {
477 TargetPath::discovery(self.target.clone())?
478 } else {
479 TargetPath::application(self.target.clone(), self.subject.clone())?
480 };
481 let mut request = http::Request::builder()
482 .method(http::Method::POST)
483 .uri(target_path.to_string())
484 .body(self.payload.clone())
485 .map_err(|source| CoreError::Malformed(source.to_string()))?;
486 if self.kind != Kind::Request {
487 request
488 .headers_mut()
489 .insert(UNB_KIND, http::HeaderValue::from_static(self.kind.token()));
490 }
491 self.project_reserved(request.headers_mut())?;
492 Self::project_custom(request.headers_mut(), &self.headers)?;
493 Ok(request)
494 }
495
496 pub fn to_local_request(&self) -> Result<http::Request<Bytes>, CoreError> {
497 if !matches!(self.kind, Kind::Request | Kind::Subscribe | Kind::Channel) {
498 return Err(CoreError::Malformed(format!(
499 "{:?} is not a destination-local handler request kind",
500 self.kind
501 )));
502 }
503 let target_path = TargetPath::application(self.target.clone(), self.subject.clone())?;
504 let mut request = http::Request::builder()
505 .method(http::Method::POST)
506 .uri(format!("/{}", target_path.subject().replace('.', "/")))
507 .body(self.payload.clone())
508 .map_err(|source| CoreError::Malformed(source.to_string()))?;
509 if self.kind != Kind::Request {
510 request
511 .headers_mut()
512 .insert(UNB_KIND, http::HeaderValue::from_static(self.kind.token()));
513 }
514 self.project_reserved(request.headers_mut())?;
515 Self::project_custom(request.headers_mut(), &self.headers)?;
516 Ok(request)
517 }
518
519 pub fn subject_of(uri: &http::Uri) -> String {
520 let path = uri.path().trim_start_matches('/');
521 if path.is_empty() {
522 uri.authority()
523 .map(|authority| authority.as_str().replace('/', "."))
524 .unwrap_or_default()
525 } else {
526 path.replace('/', ".")
527 }
528 }
529
530 pub fn ensure_headers_wire_safe(
531 headers: &serde_json::Map<String, Value>,
532 ) -> Result<(), CoreError> {
533 Self::project_custom(&mut http::HeaderMap::new(), headers)
534 }
535
536 pub fn from_request(request: http::Request<Bytes>) -> Result<Envelope, CoreError> {
537 if request.method() != http::Method::POST {
538 return Err(CoreError::Malformed(format!(
539 "application requests use POST, not {}",
540 request.method()
541 )));
542 }
543 if request.uri().query().is_some() {
544 return Err(CoreError::Malformed(
545 "application request targets must not carry a query string".into(),
546 ));
547 }
548 let kind = match request.headers().get(UNB_KIND) {
549 None => Kind::Request,
550 Some(value) => {
551 let value = value
552 .to_str()
553 .map_err(|source| CoreError::Malformed(format!("{UNB_KIND}: {source}")))?;
554 match Kind::from_token(value) {
555 Some(kind) if kind.is_application_request() => kind,
556 _ => {
557 return Err(CoreError::Malformed(format!(
558 "{UNB_KIND}: {value:?} is not an application request kind"
559 )))
560 }
561 }
562 }
563 };
564 let (parts, payload) = request.into_parts();
565 let target_path = if kind == Kind::Discover {
566 TargetPath::parse_discovery(parts.uri.path())?
567 } else {
568 TargetPath::parse_application(parts.uri.path())?
569 };
570 let mut envelope = Envelope {
571 v: PROTOCOL_VERSION,
572 id: String::new(),
573 target: target_path.target().to_owned(),
574 subject: target_path.subject().to_owned(),
575 kind,
576 corr: None,
577 seq: None,
578 hops: None,
579 body_token: None,
580 payload,
581 path: Vec::new(),
582 headers: serde_json::Map::new(),
583 };
584 envelope.absorb_headers(&parts.headers, false)?;
585 Ok(envelope)
586 }
587
588 pub fn to_response(&self) -> Result<http::Response<Bytes>, CoreError> {
589 let status = match self.kind {
590 Kind::Response | Kind::Event => http::StatusCode::OK,
591 Kind::Error => self
592 .error_code()
593 .ok_or_else(|| {
594 CoreError::Malformed("error frame payload carries no known code".into())
595 })?
596 .status(),
597 other => {
598 return Err(CoreError::Malformed(format!(
599 "{other:?} is not an application response kind"
600 )))
601 }
602 };
603 if !self.subject.is_empty() {
604 return Err(CoreError::Malformed(
605 "response frames carry no subject".into(),
606 ));
607 }
608 if self.kind == Kind::Response && self.seq.is_some() {
609 return Err(CoreError::Malformed(
610 "a unary response carries no seq; seq marks stream events".into(),
611 ));
612 }
613 if self.kind == Kind::Event && self.seq.is_none() {
614 return Err(CoreError::Malformed(
615 "an event frame requires seq to remain distinguishable".into(),
616 ));
617 }
618 let mut response = http::Response::builder()
619 .status(status)
620 .body(self.payload.clone())
621 .map_err(|source| CoreError::Malformed(source.to_string()))?;
622 self.project_reserved(response.headers_mut())?;
623 if self.kind == Kind::Error {
624 let code = self.payload_json()["code"].clone();
625 if let Some(code) = code.as_str() {
626 response.headers_mut().insert(
627 UNB_CODE,
628 http::HeaderValue::from_str(code)
629 .map_err(|source| CoreError::Malformed(source.to_string()))?,
630 );
631 }
632 }
633 Self::project_custom(response.headers_mut(), &self.headers)?;
634 Ok(response)
635 }
636
637 pub fn from_response(response: http::Response<Bytes>) -> Result<Envelope, CoreError> {
638 let status = response.status();
639 let (parts, payload) = response.into_parts();
640 let kind = if status == http::StatusCode::OK {
641 if parts.headers.contains_key(UNB_SEQ) {
642 Kind::Event
643 } else {
644 Kind::Response
645 }
646 } else if status.is_client_error() || status.is_server_error() {
647 Kind::Error
648 } else {
649 return Err(CoreError::Malformed(format!(
650 "status {status} maps to no application response kind"
651 )));
652 };
653 let mut envelope = Envelope {
654 v: PROTOCOL_VERSION,
655 id: String::new(),
656 target: String::new(),
657 subject: String::new(),
658 kind,
659 corr: None,
660 seq: None,
661 hops: None,
662 body_token: None,
663 payload,
664 path: Vec::new(),
665 headers: serde_json::Map::new(),
666 };
667 envelope.absorb_headers(&parts.headers, true)?;
668 Ok(envelope)
669 }
670
671 fn project_reserved(&self, headers: &mut http::HeaderMap) -> Result<(), CoreError> {
672 let mut put = |name: &'static str, value: String| -> Result<(), CoreError> {
673 let value = http::HeaderValue::from_str(&value)
674 .map_err(|source| CoreError::Malformed(format!("{name}: {source}")))?;
675 headers.insert(name, value);
676 Ok(())
677 };
678 put(UNB_VERSION, self.v.to_string())?;
679 if !self.id.is_empty() {
680 put(UNB_ID, self.id.clone())?;
681 }
682 if let Some(corr) = &self.corr {
683 put(UNB_CORR, corr.clone())?;
684 }
685 if let Some(seq) = self.seq {
686 put(UNB_SEQ, seq.to_string())?;
687 }
688 if let Some(hops) = self.hops {
689 put(UNB_HOPS, hops.to_string())?;
690 }
691 if !self.path.is_empty() {
692 if self.path.iter().any(|hop| hop.contains(',')) {
693 return Err(CoreError::Malformed(
694 "path elements must not contain commas".into(),
695 ));
696 }
697 put(UNB_PATH, self.path.join(","))?;
698 }
699 Ok(())
700 }
701
702 fn project_custom(
703 headers: &mut http::HeaderMap,
704 custom: &serde_json::Map<String, Value>,
705 ) -> Result<(), CoreError> {
706 for (name, value) in custom {
707 let name = http::HeaderName::try_from(name.as_str())
708 .map_err(|source| CoreError::Malformed(source.to_string()))?;
709 if name.as_str().starts_with("unb-") {
710 return Err(CoreError::Malformed(format!(
711 "{name}: unb-* header names are reserved for framing metadata"
712 )));
713 }
714 let Value::String(value) = value else {
715 return Err(CoreError::Malformed(format!(
716 "{name}: application-frame header values must be strings"
717 )));
718 };
719 let value = http::HeaderValue::from_str(value)
720 .map_err(|source| CoreError::Malformed(format!("{name}: {source}")))?;
721 headers.insert(name, value);
722 }
723 Ok(())
724 }
725
726 fn absorb_headers(
727 &mut self,
728 headers: &http::HeaderMap,
729 response: bool,
730 ) -> Result<(), CoreError> {
731 for name in headers.keys() {
732 let mut values = headers.get_all(name).iter();
733 let value = values.next().expect("keys() yields present names");
734 if values.next().is_some() {
735 return Err(CoreError::Malformed(format!(
736 "{name}: duplicate header values are not representable"
737 )));
738 }
739 let value = std::str::from_utf8(value.as_bytes())
740 .map_err(|source| CoreError::Malformed(format!("{name}: {source}")))?;
741 match name.as_str() {
742 UNB_VERSION => {
743 self.v = value
744 .parse()
745 .map_err(|_| CoreError::Malformed(format!("{name}: {value:?}")))?;
746 }
747 UNB_ID => self.id = value.to_string(),
748 UNB_CORR => self.corr = Some(value.to_string()),
749 UNB_SEQ => {
750 self.seq = Some(
751 value
752 .parse()
753 .map_err(|_| CoreError::Malformed(format!("{name}: {value:?}")))?,
754 );
755 }
756 UNB_HOPS => {
757 self.hops = Some(
758 value
759 .parse()
760 .map_err(|_| CoreError::Malformed(format!("{name}: {value:?}")))?,
761 );
762 }
763 UNB_PATH => {
764 self.path = value.split(',').map(str::to_string).collect();
765 }
766 UNB_CODE if response => {}
767 UNB_KIND if !response => {}
768 other if other.starts_with("unb-") => {
769 return Err(CoreError::Malformed(format!(
770 "{other}: unknown reserved header"
771 )));
772 }
773 other => {
774 self.headers
775 .insert(other.to_string(), Value::String(value.to_string()));
776 }
777 }
778 }
779 Ok(())
780 }
781}
782
783#[cfg(test)]
784mod tests {
785 use super::*;
786 use serde_json::json;
787
788 fn hex_bytes(hex: &str) -> Bytes {
789 let hex = hex.trim();
790 (0..hex.len())
791 .step_by(2)
792 .map(|i| u8::from_str_radix(&hex[i..i + 2], 16).expect("hex fixture"))
793 .collect::<Vec<u8>>()
794 .into()
795 }
796
797 fn envelope(kind: Kind) -> Envelope {
798 Envelope {
799 v: PROTOCOL_VERSION,
800 id: "f1".into(),
801 target: if kind.is_application_request() {
802 "node-a".into()
803 } else {
804 String::new()
805 },
806 subject: String::new(),
807 kind,
808 corr: None,
809 seq: None,
810 hops: None,
811 body_token: None,
812 payload: Bytes::new(),
813 path: Vec::new(),
814 headers: Default::default(),
815 }
816 }
817
818 #[test]
819 fn a_v1_request_line_frame_round_trips() {
820 let mut headers = serde_json::Map::new();
821 headers.insert("authorization".into(), json!("Bearer jwt-abc"));
822 let request = Envelope {
823 v: PROTOCOL_VERSION,
824 id: "f2".into(),
825 target: "node-b".into(),
826 subject: "chess.move".into(),
827 kind: Kind::Request,
828 corr: Some("s1".into()),
829 seq: None,
830 hops: Some(DEFAULT_HOPS),
831 body_token: None,
832 payload: Envelope::encode_payload(&json!({"from": "e2", "to": "e4"})),
833 path: vec!["node-a".into()],
834 headers,
835 };
836 let frame = request.encode();
837 let text = std::str::from_utf8(&frame).unwrap();
838 assert_eq!(
839 text,
840 "POST /node-b/chess/move HTTP/1.1\r\nunb-version: 1\r\nunb-id: f2\r\nunb-corr: s1\r\nunb-hops: 8\r\n\
841 unb-path: node-a\r\nauthorization: Bearer jwt-abc\r\n\r\n\
842 {\"from\":\"e2\",\"to\":\"e4\"}"
843 );
844 assert_eq!(Envelope::decode(frame).unwrap(), request);
845 }
846
847 #[test]
848 fn every_kind_takes_its_lane_on_the_wire() {
849 for kind in [
850 Kind::Hello,
851 Kind::Welcome,
852 Kind::Request,
853 Kind::Response,
854 Kind::Subscribe,
855 Kind::Event,
856 Kind::Channel,
857 Kind::Cancel,
858 Kind::Error,
859 Kind::Ping,
860 Kind::Pong,
861 Kind::Discover,
862 Kind::Identify,
863 Kind::IdentityAccepted,
864 Kind::RouteSnapshot,
865 Kind::RouteDelta,
866 Kind::RouteAck,
867 ] {
868 let mut envelope = envelope(kind);
869 if kind.is_application_request() && kind != Kind::Discover {
870 envelope.subject = "chess".into();
871 envelope.corr = Some("s1".into());
872 }
873 if kind == Kind::Discover {
874 envelope.corr = Some("s1".into());
875 }
876 if kind.is_application_response() {
877 envelope.corr = Some("s1".into());
878 }
879 if kind == Kind::Event {
880 envelope.seq = Some(7);
881 }
882 if kind == Kind::Error {
883 envelope.payload =
884 Envelope::encode_payload(&json!({"code": "BUSY", "message": "full"}));
885 }
886 let frame = envelope.encode();
887 let text = std::str::from_utf8(&frame).unwrap();
888 if kind.is_application_response() {
889 assert!(text.starts_with("HTTP/1.1 "), "{kind:?}: {text}");
890 } else if kind == Kind::Discover {
891 assert!(
892 text.starts_with("POST /node-a HTTP/1.1"),
893 "{kind:?}: {text}"
894 );
895 } else if kind.is_application_request() {
896 assert!(
897 text.starts_with("POST /node-a/chess HTTP/1.1"),
898 "{kind:?}: {text}"
899 );
900 } else {
901 assert!(text.starts_with("POST * HTTP/1.1"), "{kind:?}: {text}");
902 }
903 assert_eq!(Envelope::decode(frame).unwrap(), envelope, "{kind:?}");
904 }
905 }
906
907 #[test]
908 fn a_control_frame_routes_by_az_kind_over_the_asterisk_target() {
909 let mut hello = envelope(Kind::Hello);
910 hello.payload = Envelope::encode_payload(&json!({"versions": [1]}));
911 assert_eq!(
912 std::str::from_utf8(&hello.encode()).unwrap(),
913 "POST * HTTP/1.1\r\nunb-version: 1\r\nunb-id: f1\r\nunb-kind: hello\r\n\r\n{\"versions\":[1]}"
914 );
915 }
916
917 #[test]
918 fn a_ping_frame_has_an_empty_body() {
919 let ping = envelope(Kind::Ping);
920 let frame = ping.encode();
921 assert_eq!(
922 std::str::from_utf8(&frame).unwrap(),
923 "POST * HTTP/1.1\r\nunb-version: 1\r\nunb-id: f1\r\nunb-kind: ping\r\n\r\n"
924 );
925 assert_eq!(Envelope::decode(frame).unwrap(), ping);
926 }
927
928 #[test]
929 fn an_error_status_line_names_the_error_code() {
930 let mut error = envelope(Kind::Error);
931 error.corr = Some("s1".into());
932 error.payload =
933 Envelope::encode_payload(&json!({"code": "BUSY", "message": "node at capacity"}));
934 let frame = error.encode();
935 let text = std::str::from_utf8(&frame).unwrap();
936 assert!(text.starts_with("HTTP/1.1 503 BUSY\r\n"), "{text}");
937 assert!(text.contains("\r\nunb-code: BUSY\r\n"), "{text}");
938 assert_eq!(Envelope::decode(frame).unwrap(), error);
939 }
940
941 #[test]
942 fn every_error_code_has_canonical_wire_status_reason_and_restoration() {
943 for (code, status) in [
944 (crate::ErrorCode::VersionMismatch, 505),
945 (crate::ErrorCode::Protocol, 400),
946 (crate::ErrorCode::UnknownSubject, 404),
947 (crate::ErrorCode::UnresolvedAtPeer, 421),
948 (crate::ErrorCode::PeerUnreachable, 502),
949 (crate::ErrorCode::HopLimitExceeded, 508),
950 (crate::ErrorCode::InvalidInput, 400),
951 (crate::ErrorCode::Unauthorized, 401),
952 (crate::ErrorCode::Conflict, 409),
953 (crate::ErrorCode::Busy, 503),
954 (crate::ErrorCode::Cancelled, 499),
955 (crate::ErrorCode::Internal, 500),
956 ] {
957 let mut error = envelope(Kind::Error);
958 error.corr = Some("s1".into());
959 error.payload = Envelope::encode_payload(&json!({
960 "code": code.token(),
961 "message": "failure"
962 }));
963 let frame = error.encode();
964 let text = std::str::from_utf8(&frame).unwrap();
965 assert!(
966 text.starts_with(&format!("HTTP/1.1 {status} {}\r\n", code.token())),
967 "{code:?}: {text}"
968 );
969 assert!(
970 text.contains(&format!("\r\nunb-code: {}\r\n", code.token())),
971 "{code:?}: {text}"
972 );
973 let decoded = Envelope::decode(frame.clone()).unwrap();
974 assert_eq!(decoded, error, "{code:?}");
975 let restored = Envelope::from_response(decoded.to_response().unwrap()).unwrap();
976 assert_eq!(restored.payload_json()["code"], code.token(), "{code:?}");
977 assert_eq!(restored.encode(), frame, "{code:?}");
978 }
979 }
980
981 #[test]
982 fn an_unmapped_error_payload_still_frames_as_a_server_error() {
983 let mut error = envelope(Kind::Error);
984 error.corr = Some("s1".into());
985 error.payload = Envelope::encode_payload(&json!({"note": "no code field"}));
986 let frame = error.encode();
987 let text = std::str::from_utf8(&frame).unwrap();
988 assert!(text.starts_with("HTTP/1.1 500 INTERNAL\r\n"), "{text}");
989 let decoded = Envelope::decode(frame).unwrap();
990 assert_eq!(decoded.kind, Kind::Error);
991 }
992
993 #[test]
994 fn decode_slices_the_body_zero_copy() {
995 let mut request = envelope(Kind::Request);
996 request.subject = "chess".into();
997 request.corr = Some("s1".into());
998 request.payload = Envelope::encode_payload(&json!({"from": "e2"}));
999 let frame = request.encode();
1000 let body_start = frame.len() - request.payload.len();
1001 let decoded = Envelope::decode(frame.clone()).unwrap();
1002 assert_eq!(decoded.payload.as_ptr(), frame[body_start..].as_ptr());
1003 }
1004
1005 #[test]
1006 fn dot_and_slash_subject_forms_are_equivalent() {
1007 for uri in ["/node-a/chess.move", "/node-a/chess/move"] {
1008 let request = http::Request::builder()
1009 .method("POST")
1010 .uri(uri)
1011 .body(Bytes::new())
1012 .unwrap();
1013 let envelope = Envelope::from_request(request).unwrap();
1014 assert_eq!(envelope.target, "node-a");
1015 assert_eq!(envelope.subject, "chess.move");
1016 }
1017 }
1018
1019 #[test]
1020 fn a_malformed_head_is_a_typed_decode_error() {
1021 for frame in [
1022 &b"GARBAGE\r\n\r\n"[..],
1023 b"POST /x HTTP/1.1\r\nunb-id: f1",
1024 b"GET /x HTTP/1.1\r\n\r\n",
1025 b"POST x HTTP/1.1\r\n\r\n",
1026 b"POST /x?side=w HTTP/1.1\r\nunb-version: 1\r\n\r\n",
1027 b"POST /az/teleport HTTP/1.1\r\nunb-version: 1\r\n\r\n",
1028 b"POST /az HTTP/1.1\r\nunb-version: 1\r\n\r\n",
1029 b"POST /az/move HTTP/1.1\r\nunb-version: 1\r\n\r\n",
1030 b"POST /x HTTP/1.1\r\nunb-kind: response\r\n\r\n",
1031 b"POST /x HTTP/1.1\r\nunb-magic: 1\r\n\r\n",
1032 b"POST /x HTTP/1.1\r\nunb-corr: a\r\nunb-corr: b\r\n\r\n",
1033 b"POST /x HTTP/1.1\r\nactor: a\r\nactor: b\r\n\r\n",
1034 b"POST * HTTP/1.1\r\nunb-kind: request\r\n\r\n",
1035 b"POST * HTTP/1.1\r\n\r\n",
1036 b"POST /x HTTP/1.1\r\ncontent-length: 5\r\n\r\nhello",
1037 b"POST /x HTTP/1.1\r\ntransfer-encoding: chunked\r\n\r\n",
1038 b"HTTP/1.1 302 FOUND\r\n\r\n",
1039 b"HTTP/1.1 200 OK\r\nunb-kind: response\r\n\r\n",
1040 ] {
1041 let error = Envelope::decode(Bytes::copy_from_slice(frame)).unwrap_err();
1042 assert!(
1043 matches!(error, CoreError::Malformed(_)),
1044 "{:?}: {error}",
1045 std::str::from_utf8(frame)
1046 );
1047 }
1048 }
1049
1050 #[test]
1051 fn decode_rejects_bare_carriage_returns_and_accepts_supported_line_endings() {
1052 for frame in [
1053 &b"POST /x\rHTTP/1.1\r\n\r\n"[..],
1054 &b"POST /x HTTP/1.1\r\nact\ror: a\r\n\r\n"[..],
1055 &b"POST /x HTTP/1.1\r\nactor: a\rb\r\n\r\n"[..],
1056 ] {
1057 assert!(matches!(
1058 Envelope::decode(Bytes::copy_from_slice(frame)),
1059 Err(CoreError::Malformed(_))
1060 ));
1061 }
1062
1063 assert_eq!(
1064 Envelope::decode(Bytes::from_static(
1065 b"POST /node-a/x HTTP/1.1\r\nactor: a\r\n\r\n",
1066 ))
1067 .unwrap()
1068 .headers["actor"],
1069 "a"
1070 );
1071 }
1072
1073 #[test]
1074 fn the_exact_head_bound_is_accepted_and_one_byte_over_is_rejected() {
1075 let prefix = "POST /node-a/x HTTP/1.1\r\nactor: ";
1076 let suffix = "\r\n\r\n";
1077 let value_len = MAX_HEAD_BYTES - prefix.len() - suffix.len();
1078 let exact = format!("{prefix}{}{suffix}", "x".repeat(value_len));
1079 assert_eq!(exact.len(), MAX_HEAD_BYTES);
1080 assert_eq!(
1081 Envelope::decode(Bytes::from(exact)).unwrap().headers["actor"],
1082 "x".repeat(value_len)
1083 );
1084
1085 let text = format!("{prefix}{}{suffix}", "x".repeat(value_len + 1));
1086 assert_eq!(text.len(), MAX_HEAD_BYTES + 1);
1087 assert!(matches!(
1088 Envelope::decode(Bytes::from(text)),
1089 Err(CoreError::Malformed(_))
1090 ));
1091 }
1092
1093 #[test]
1094 fn custom_headers_stay_headers_not_core_fields() {
1095 let frame = Bytes::from_static(
1096 b"POST /node-a/chess HTTP/1.1\r\nunb-corr: s1\r\nsubject: sneaky\r\n\r\n",
1097 );
1098 let decoded = Envelope::decode(frame).unwrap();
1099 assert_eq!(decoded.target, "node-a");
1100 assert_eq!(decoded.subject, "chess");
1101 assert_eq!(decoded.corr.as_deref(), Some("s1"));
1102 assert_eq!(decoded.headers["subject"], "sneaky");
1103 }
1104
1105 #[test]
1106 fn a_non_json_payload_round_trips_verbatim() {
1107 let mut opaque = envelope(Kind::Request);
1108 opaque.subject = "files".into();
1109 opaque.corr = Some("s1".into());
1110 opaque.payload = Bytes::from_static(&[0x00, 0x01, 0xff, 0xfe, b'!', 0x80]);
1111 let decoded = Envelope::decode(opaque.encode()).unwrap();
1112 assert_eq!(decoded, opaque);
1113 assert_eq!(decoded.payload_json(), Value::Null);
1114 }
1115
1116 #[test]
1117 fn a_custom_header_round_trips_verbatim() {
1118 let mut headers = serde_json::Map::new();
1119 headers.insert("actor".into(), json!("jwt-abc"));
1120 headers.insert("x-trace".into(), json!("span-7"));
1121 let mut request = envelope(Kind::Request);
1122 request.subject = "chess".into();
1123 request.corr = Some("s1".into());
1124 request.headers = headers;
1125 let decoded = Envelope::decode(request.encode()).unwrap();
1126 assert_eq!(decoded, request);
1127 assert_eq!(decoded.headers["actor"], "jwt-abc");
1128 }
1129
1130 #[test]
1131 fn a_mixed_case_unb_header_is_rejected_as_reserved() {
1132 for name in ["Unb-Corr", "UNB-KIND", "uNb-hOpS"] {
1133 let mut headers = serde_json::Map::new();
1134 headers.insert(name.into(), json!("forged"));
1135 assert!(matches!(
1136 Envelope::ensure_headers_wire_safe(&headers),
1137 Err(CoreError::Malformed(message)) if message.contains("reserved")
1138 ));
1139 let mut request = envelope(Kind::Request);
1140 request.subject = "chess".into();
1141 request.corr = Some("s1".into());
1142 request.headers.insert(name.into(), json!("forged"));
1143 assert!(matches!(
1144 request.to_request(),
1145 Err(CoreError::Malformed(message)) if message.contains("reserved")
1146 ));
1147 }
1148 }
1149
1150 #[test]
1151 fn an_application_envelope_round_trips_through_the_request_model() {
1152 let mut headers = serde_json::Map::new();
1153 headers.insert("authorization".into(), json!("Bearer jwt-abc"));
1154 let envelope = Envelope {
1155 v: PROTOCOL_VERSION,
1156 id: "f7".into(),
1157 target: "node-a".into(),
1158 subject: "chess.move".into(),
1159 kind: Kind::Request,
1160 corr: Some("s1".into()),
1161 seq: None,
1162 hops: Some(DEFAULT_HOPS),
1163 body_token: None,
1164 payload: Envelope::encode_payload(&json!({"from": "e2", "to": "e4"})),
1165 path: vec!["node-a".into()],
1166 headers,
1167 };
1168 let request = envelope.to_request().unwrap();
1169 assert_eq!(request.method(), http::Method::POST);
1170 assert_eq!(request.uri().path(), "/node-a/chess/move");
1171 assert_eq!(request.headers()["authorization"], "Bearer jwt-abc");
1172 assert_eq!(request.headers()[UNB_CORR], "s1");
1173 assert_eq!(request.headers()[UNB_ID], "f7");
1174 let back = Envelope::from_request(request).unwrap();
1175 assert_eq!(back.encode(), envelope.encode());
1176 }
1177
1178 #[test]
1179 fn every_request_kind_rides_az_kind_over_post_and_back() {
1180 for (kind, marker) in [
1181 (Kind::Request, None),
1182 (Kind::Subscribe, Some("subscribe")),
1183 (Kind::Channel, Some("channel")),
1184 (Kind::Discover, Some("discover")),
1185 ] {
1186 let mut envelope = envelope(kind);
1187 if kind != Kind::Discover {
1188 envelope.subject = "chess".into();
1189 }
1190 envelope.corr = Some("s1".into());
1191 let request = envelope.to_request().unwrap();
1192 assert_eq!(request.method(), http::Method::POST);
1193 match marker {
1194 Some(marker) => assert_eq!(request.headers()[UNB_KIND], marker),
1195 None => assert!(!request.headers().contains_key(UNB_KIND)),
1196 }
1197 assert_eq!(Envelope::from_request(request).unwrap().kind, kind);
1198 }
1199 }
1200
1201 #[test]
1202 fn request_conversion_rejects_what_the_model_forbids() {
1203 let mut base = envelope(Kind::Request);
1204 base.subject = "chess".into();
1205 base.corr = Some("s1".into());
1206
1207 let mut control = base.clone();
1208 control.kind = Kind::Ping;
1209 assert!(matches!(control.to_request(), Err(CoreError::Malformed(_))));
1210
1211 let mut nested = base.clone();
1212 nested.headers.insert("trace".into(), json!({"span": 7}));
1213 assert!(matches!(nested.to_request(), Err(CoreError::Malformed(_))));
1214
1215 let mut shadowing = base.clone();
1216 shadowing.headers.insert(UNB_CORR.into(), json!("spoof"));
1217 assert!(matches!(
1218 shadowing.to_request(),
1219 Err(CoreError::Malformed(_))
1220 ));
1221
1222 let mut queried = base.clone();
1223 queried.subject = "chess.move?side=white".into();
1224 assert!(matches!(queried.to_request(), Err(CoreError::Malformed(_))));
1225
1226 let mut reserved = base.clone();
1227 reserved.target = "az".into();
1228 assert!(matches!(
1229 reserved.to_request(),
1230 Err(CoreError::Malformed(_))
1231 ));
1232
1233 let wrong_method = http::Request::builder()
1234 .method("GET")
1235 .uri("/node-a/chess")
1236 .body(Bytes::new())
1237 .unwrap();
1238 assert!(matches!(
1239 Envelope::from_request(wrong_method),
1240 Err(CoreError::Malformed(_))
1241 ));
1242
1243 let unknown_kind = http::Request::builder()
1244 .method("POST")
1245 .uri("/node-a/chess")
1246 .header(UNB_KIND, "teleport")
1247 .body(Bytes::new())
1248 .unwrap();
1249 assert!(matches!(
1250 Envelope::from_request(unknown_kind),
1251 Err(CoreError::Malformed(_))
1252 ));
1253
1254 let queried = http::Request::builder()
1255 .method("POST")
1256 .uri("/node-a/chess/move?draft=1")
1257 .body(Bytes::new())
1258 .unwrap();
1259 assert!(matches!(
1260 Envelope::from_request(queried),
1261 Err(CoreError::Malformed(message)) if message.contains("query string")
1262 ));
1263
1264 let reserved_subject = http::Request::builder()
1265 .method("POST")
1266 .uri("/az/hello")
1267 .body(Bytes::new())
1268 .unwrap();
1269 assert!(matches!(
1270 Envelope::from_request(reserved_subject),
1271 Err(CoreError::Malformed(_))
1272 ));
1273 }
1274
1275 #[test]
1276 fn a_fresh_user_request_converts_with_wire_defaults() {
1277 let request = http::Request::builder()
1278 .method("POST")
1279 .uri("/node-a/chess/move")
1280 .header("authorization", "jwt-abc")
1281 .body(Bytes::from_static(b"{}"))
1282 .unwrap();
1283 let envelope = Envelope::from_request(request).unwrap();
1284 assert_eq!(envelope.v, PROTOCOL_VERSION);
1285 assert_eq!(envelope.kind, Kind::Request);
1286 assert_eq!(envelope.target, "node-a");
1287 assert_eq!(envelope.subject, "chess.move");
1288 assert!(envelope.id.is_empty());
1289 assert_eq!(envelope.corr, None);
1290 assert_eq!(envelope.headers["authorization"], "jwt-abc");
1291 }
1292
1293 #[test]
1294 fn model_conversion_shares_the_payload_allocation() {
1295 let mut source = envelope(Kind::Request);
1296 source.subject = "chess.move".into();
1297 source.corr = Some("s1".into());
1298 source.payload = Envelope::encode_payload(&json!({"from": "e2"}));
1299 let request = source.to_request().unwrap();
1300 assert_eq!(request.body().as_ptr(), source.payload.as_ptr());
1301 let back = Envelope::from_request(request).unwrap();
1302 assert_eq!(back.payload.as_ptr(), source.payload.as_ptr());
1303 }
1304
1305 #[test]
1306 fn response_event_and_error_envelopes_round_trip_through_the_response_model() {
1307 let mut response = envelope(Kind::Response);
1308 response.id = "f8".into();
1309 response.corr = Some("s1".into());
1310 response.payload = Envelope::encode_payload(&json!({"ok": true}));
1311 let converted = response.to_response().unwrap();
1312 assert_eq!(converted.status(), http::StatusCode::OK);
1313 let back = Envelope::from_response(converted).unwrap();
1314 assert_eq!(back.encode(), response.encode());
1315
1316 let mut event = response.clone();
1317 event.kind = Kind::Event;
1318 event.seq = Some(42);
1319 let converted = event.to_response().unwrap();
1320 assert_eq!(converted.headers()[UNB_SEQ], "42");
1321 let back = Envelope::from_response(converted).unwrap();
1322 assert_eq!(back.kind, Kind::Event);
1323 assert_eq!(back.encode(), event.encode());
1324
1325 let mut error = response.clone();
1326 error.kind = Kind::Error;
1327 error.payload =
1328 Envelope::encode_payload(&json!({"code": "BUSY", "message": "node at capacity"}));
1329 let converted = error.to_response().unwrap();
1330 assert_eq!(converted.status(), http::StatusCode::SERVICE_UNAVAILABLE);
1331 assert_eq!(converted.headers()[UNB_CODE], "BUSY");
1332 let back = Envelope::from_response(converted).unwrap();
1333 assert_eq!(back.encode(), error.encode());
1334 }
1335
1336 #[test]
1337 fn reserved_fields_project_and_restore_across_legal_application_lanes() {
1338 for kind in [Kind::Request, Kind::Event, Kind::Error] {
1339 let mut source = envelope(kind);
1340 source.id = "f-reserved".into();
1341 source.subject = if kind == Kind::Request {
1342 "chess.move".into()
1343 } else {
1344 String::new()
1345 };
1346 source.corr = Some("s-reserved".into());
1347 source.seq = (kind == Kind::Event).then_some(17);
1348 source.hops = Some(6);
1349 source.path = vec!["leaf-a".into(), "hub-b".into(), "root-c".into()];
1350 source.headers.insert("x-trace".into(), json!("span-7"));
1351 source.payload = if kind == Kind::Error {
1352 Envelope::encode_payload(&json!({"code": "PROTOCOL", "message": "bad"}))
1353 } else {
1354 Bytes::from_static(b"opaque")
1355 };
1356
1357 let restored = if kind == Kind::Request {
1358 let projected = source.to_request().unwrap();
1359 assert_eq!(projected.headers()[UNB_VERSION], "1");
1360 assert_eq!(projected.headers()[UNB_ID], "f-reserved");
1361 assert_eq!(projected.headers()[UNB_CORR], "s-reserved");
1362 assert_eq!(projected.headers()[UNB_HOPS], "6");
1363 assert_eq!(projected.headers()[UNB_PATH], "leaf-a,hub-b,root-c");
1364 assert_eq!(projected.headers()["x-trace"], "span-7");
1365 Envelope::from_request(projected).unwrap()
1366 } else {
1367 let projected = source.to_response().unwrap();
1368 assert_eq!(projected.headers()[UNB_VERSION], "1");
1369 assert_eq!(projected.headers()[UNB_ID], "f-reserved");
1370 assert_eq!(projected.headers()[UNB_CORR], "s-reserved");
1371 assert_eq!(projected.headers()[UNB_HOPS], "6");
1372 assert_eq!(projected.headers()[UNB_PATH], "leaf-a,hub-b,root-c");
1373 assert_eq!(projected.headers()["x-trace"], "span-7");
1374 if kind == Kind::Event {
1375 assert_eq!(projected.headers()[UNB_SEQ], "17");
1376 } else {
1377 assert_eq!(projected.headers()[UNB_CODE], "PROTOCOL");
1378 }
1379 Envelope::from_response(projected).unwrap()
1380 };
1381 assert_eq!(restored, source, "{kind:?}");
1382 }
1383 }
1384
1385 #[test]
1386 fn comma_path_elements_are_rejected_on_every_legal_projection_lane() {
1387 for kind in [Kind::Request, Kind::Event, Kind::Error] {
1388 let mut source = envelope(kind);
1389 source.subject = if kind == Kind::Request {
1390 "chess".into()
1391 } else {
1392 String::new()
1393 };
1394 source.corr = Some("s1".into());
1395 source.seq = (kind == Kind::Event).then_some(1);
1396 source.path = vec!["leaf-a,forged-hop".into()];
1397 if kind == Kind::Error {
1398 source.payload = Envelope::encode_payload(&json!({"code": "INTERNAL"}));
1399 }
1400 let result = if kind == Kind::Request {
1401 source.to_request().map(|_| ())
1402 } else {
1403 source.to_response().map(|_| ())
1404 };
1405 assert!(matches!(result, Err(CoreError::Malformed(_))), "{kind:?}");
1406 }
1407 }
1408
1409 #[test]
1410 fn error_status_collisions_restore_the_exact_code_from_az_code() {
1411 for code in ["INVALID_INPUT", "PROTOCOL"] {
1412 let mut error = envelope(Kind::Error);
1413 error.id = "f9".into();
1414 error.corr = Some("s1".into());
1415 error.payload = Envelope::encode_payload(&json!({"code": code, "message": "bad"}));
1416 let converted = error.to_response().unwrap();
1417 assert_eq!(converted.status(), http::StatusCode::BAD_REQUEST);
1418 assert_eq!(converted.headers()[UNB_CODE], code);
1419 let back = Envelope::from_response(converted).unwrap();
1420 assert_eq!(back.encode(), error.encode());
1421 }
1422 }
1423
1424 #[test]
1425 fn header_validation_rejects_frame_splitting_input() {
1426 let clean = {
1427 let mut headers = serde_json::Map::new();
1428 headers.insert("x-trace".into(), json!("span-7"));
1429 headers
1430 };
1431 assert!(Envelope::ensure_headers_wire_safe(&clean).is_ok());
1432
1433 let crlf_value = {
1434 let mut headers = serde_json::Map::new();
1435 headers.insert("x-trace".into(), json!("span\r\nunb-corr: forged"));
1436 headers
1437 };
1438 assert!(matches!(
1439 Envelope::ensure_headers_wire_safe(&crlf_value),
1440 Err(CoreError::Malformed(_))
1441 ));
1442
1443 let crlf_name = {
1444 let mut headers = serde_json::Map::new();
1445 headers.insert("x\r\nInjected".into(), json!("v"));
1446 headers
1447 };
1448 assert!(matches!(
1449 Envelope::ensure_headers_wire_safe(&crlf_name),
1450 Err(CoreError::Malformed(_))
1451 ));
1452
1453 let non_string = {
1454 let mut headers = serde_json::Map::new();
1455 headers.insert("x-trace".into(), json!({ "nested": true }));
1456 headers
1457 };
1458 assert!(matches!(
1459 Envelope::ensure_headers_wire_safe(&non_string),
1460 Err(CoreError::Malformed(_))
1461 ));
1462 }
1463
1464 fn golden_envelopes() -> Vec<(&'static str, Envelope)> {
1465 let build = |id: &str, subject: &str, kind: Kind, corr: Option<&str>| Envelope {
1466 v: PROTOCOL_VERSION,
1467 id: id.into(),
1468 target: if kind.is_application_request() {
1469 "node-a".into()
1470 } else {
1471 String::new()
1472 },
1473 subject: subject.into(),
1474 kind,
1475 corr: corr.map(str::to_string),
1476 seq: None,
1477 hops: None,
1478 body_token: None,
1479 payload: Bytes::new(),
1480 path: Vec::new(),
1481 headers: Default::default(),
1482 };
1483 let mut request = build("f2", "chess", Kind::Request, Some("s1"));
1484 request.hops = Some(DEFAULT_HOPS);
1485 request.payload = Envelope::encode_payload(
1486 &json!({"action": "move", "input": {"from": "e2", "to": "e4"}}),
1487 );
1488 let mut event = build("f9", "", Kind::Event, Some("s1"));
1489 event.seq = Some(3);
1490 event.payload =
1491 Envelope::encode_payload(&json!({"as_of": 4711, "value": {"done": false, "id": "t1"}}));
1492 let mut response = build("f3", "", Kind::Response, Some("s1"));
1493 response.payload = Envelope::encode_payload(&json!({"ok": true}));
1494 let mut subscribe = build("f8", "todo.changes", Kind::Subscribe, Some("s2"));
1495 subscribe.hops = Some(DEFAULT_HOPS);
1496 subscribe.payload = Envelope::encode_payload(&json!({"after": 4711}));
1497 let mut error = build("f4", "", Kind::Error, Some("s1"));
1498 error.path = vec!["node-a".into(), "node-b".into()];
1499 error.payload = Envelope::encode_payload(&json!({
1500 "code": "UNKNOWN_SUBJECT",
1501 "message": "Unknown subject \"ches\". Did you mean \"chess\"?"
1502 }));
1503 let ping = build("f1", "", Kind::Ping, None);
1504 let mut identify = build("f1", "", Kind::Identify, None);
1505 identify.payload = Envelope::encode_payload(
1506 &json!({"node_id": "node-a", "instance_id": "node-a-1", "epoch": 1}),
1507 );
1508 let identity_accepted = build("f2", "", Kind::IdentityAccepted, None);
1509 let mut route_snapshot = build("f3", "", Kind::RouteSnapshot, None);
1510 route_snapshot.payload = Envelope::encode_payload(&json!({
1511 "generation": 1,
1512 "routes": [{
1513 "subject": "chess", "owner": "leaf-a", "owner_instance": "leaf-a-1",
1514 "owner_epoch": 1, "owner_revision": 0, "distance": 1,
1515 "path": ["leaf-a", "hub"]
1516 }]
1517 }));
1518 let mut route_delta = build("f4", "", Kind::RouteDelta, None);
1519 route_delta.payload = Envelope::encode_payload(&json!({
1520 "generation": 2,
1521 "upsert": [],
1522 "withdraw": [{
1523 "subject": "chess", "owner": "leaf-a", "owner_instance": "leaf-a-1",
1524 "owner_epoch": 1, "owner_revision": 0
1525 }]
1526 }));
1527 let mut route_ack = build("f5", "", Kind::RouteAck, None);
1528 route_ack.payload =
1529 Envelope::encode_payload(&json!({"generation": 2, "status": "applied"}));
1530 vec![
1531 ("request", request),
1532 ("response", response),
1533 ("subscribe", subscribe),
1534 ("event", event),
1535 ("error", error),
1536 ("ping", ping),
1537 ("identify", identify),
1538 ("identity_accepted", identity_accepted),
1539 ("route_snapshot", route_snapshot),
1540 ("route_delta", route_delta),
1541 ("route_ack", route_ack),
1542 ]
1543 }
1544
1545 #[test]
1546 #[ignore = "regenerates the golden fixtures from the canonical encoder"]
1547 fn regenerate_golden_fixtures() {
1548 let root = concat!(
1549 env!("CARGO_MANIFEST_DIR"),
1550 "/../../conformance/fixtures/envelope"
1551 );
1552 for (name, envelope) in golden_envelopes() {
1553 let frame = envelope.encode();
1554 let hex: String = frame.iter().map(|byte| format!("{byte:02x}")).collect();
1555 std::fs::write(format!("{root}/{name}.frame.hex"), format!("{hex}\n")).unwrap();
1556 std::fs::write(
1557 format!("{root}/{name}.header.json"),
1558 format!("{}\n", serde_json::to_string(&envelope).unwrap()),
1559 )
1560 .unwrap();
1561 }
1562 }
1563
1564 #[test]
1565 fn az_body_is_a_core_field_on_the_wire_and_reserved_for_users() {
1566 let mut headers = serde_json::Map::new();
1567 headers.insert(UNB_BODY.into(), Value::String("token".into()));
1568 assert!(matches!(
1569 Envelope::ensure_headers_wire_safe(&headers),
1570 Err(CoreError::Malformed(_))
1571 ));
1572 let frame = Bytes::from_static(
1573 b"POST /node-a/echo HTTP/1.1\r\nunb-id: a\r\nunb-body: token\r\n\r\n",
1574 );
1575 let decoded = Envelope::decode(frame).unwrap();
1576 assert_eq!(decoded.body_token.as_deref(), Some("token"));
1577 assert!(!decoded.headers.contains_key(UNB_BODY));
1578 let round = Envelope::decode(decoded.encode()).unwrap();
1579 assert_eq!(round.body_token.as_deref(), Some("token"));
1580 let oversized = format!(
1581 "POST /node-a/echo HTTP/1.1\r\nunb-id: a\r\nunb-body: {}\r\n\r\n",
1582 "x".repeat(257)
1583 );
1584 assert!(matches!(
1585 Envelope::decode(Bytes::from(oversized)),
1586 Err(CoreError::Malformed(_))
1587 ));
1588 }
1589
1590 #[test]
1591 fn head_first_decode_matches_whole_message_decode() {
1592 for (name, envelope) in golden_envelopes() {
1593 let frame = envelope.encode();
1594 let whole = Envelope::decode(frame.clone()).unwrap();
1595 let (mut head, head_len) = Envelope::decode_head(&frame).unwrap();
1596 assert!(
1597 head.payload.is_empty(),
1598 "{name}: head decode carries no payload"
1599 );
1600 head.payload = frame.slice(head_len..);
1601 assert_eq!(head, whole, "{name}");
1602 }
1603 }
1604
1605 #[test]
1606 fn golden_fixtures_round_trip_byte_identically() {
1607 for (frame_hex, header_json) in [
1608 (
1609 include_str!("../../../conformance/fixtures/envelope/request.frame.hex"),
1610 include_str!("../../../conformance/fixtures/envelope/request.header.json"),
1611 ),
1612 (
1613 include_str!("../../../conformance/fixtures/envelope/response.frame.hex"),
1614 include_str!("../../../conformance/fixtures/envelope/response.header.json"),
1615 ),
1616 (
1617 include_str!("../../../conformance/fixtures/envelope/subscribe.frame.hex"),
1618 include_str!("../../../conformance/fixtures/envelope/subscribe.header.json"),
1619 ),
1620 (
1621 include_str!("../../../conformance/fixtures/envelope/event.frame.hex"),
1622 include_str!("../../../conformance/fixtures/envelope/event.header.json"),
1623 ),
1624 (
1625 include_str!("../../../conformance/fixtures/envelope/error.frame.hex"),
1626 include_str!("../../../conformance/fixtures/envelope/error.header.json"),
1627 ),
1628 (
1629 include_str!("../../../conformance/fixtures/envelope/ping.frame.hex"),
1630 include_str!("../../../conformance/fixtures/envelope/ping.header.json"),
1631 ),
1632 (
1633 include_str!("../../../conformance/fixtures/envelope/identify.frame.hex"),
1634 include_str!("../../../conformance/fixtures/envelope/identify.header.json"),
1635 ),
1636 (
1637 include_str!("../../../conformance/fixtures/envelope/identity_accepted.frame.hex"),
1638 include_str!(
1639 "../../../conformance/fixtures/envelope/identity_accepted.header.json"
1640 ),
1641 ),
1642 (
1643 include_str!("../../../conformance/fixtures/envelope/route_snapshot.frame.hex"),
1644 include_str!("../../../conformance/fixtures/envelope/route_snapshot.header.json"),
1645 ),
1646 (
1647 include_str!("../../../conformance/fixtures/envelope/route_delta.frame.hex"),
1648 include_str!("../../../conformance/fixtures/envelope/route_delta.header.json"),
1649 ),
1650 (
1651 include_str!("../../../conformance/fixtures/envelope/route_ack.frame.hex"),
1652 include_str!("../../../conformance/fixtures/envelope/route_ack.header.json"),
1653 ),
1654 ] {
1655 let frame = hex_bytes(frame_hex);
1656 let envelope = Envelope::decode(frame.clone()).unwrap();
1657 assert_eq!(
1658 serde_json::to_string(&envelope).unwrap(),
1659 header_json.trim(),
1660 "header fixture must be the canonical envelope encoding"
1661 );
1662 assert_eq!(
1663 envelope.encode(),
1664 frame,
1665 "fixture must be the canonical frame encoding"
1666 );
1667 }
1668 }
1669
1670 #[test]
1671 fn application_golden_fixtures_survive_the_model_byte_identically() {
1672 let request_frame = hex_bytes(include_str!(
1673 "../../../conformance/fixtures/envelope/request.frame.hex"
1674 ));
1675 let envelope = Envelope::decode(request_frame.clone()).unwrap();
1676 let back = Envelope::from_request(envelope.to_request().unwrap()).unwrap();
1677 assert_eq!(back.encode(), request_frame);
1678
1679 let subscribe_frame = hex_bytes(include_str!(
1680 "../../../conformance/fixtures/envelope/subscribe.frame.hex"
1681 ));
1682 let envelope = Envelope::decode(subscribe_frame.clone()).unwrap();
1683 let back = Envelope::from_request(envelope.to_request().unwrap()).unwrap();
1684 assert_eq!(back.encode(), subscribe_frame);
1685
1686 for fixture in [
1687 include_str!("../../../conformance/fixtures/envelope/response.frame.hex"),
1688 include_str!("../../../conformance/fixtures/envelope/event.frame.hex"),
1689 include_str!("../../../conformance/fixtures/envelope/error.frame.hex"),
1690 ] {
1691 let frame = hex_bytes(fixture);
1692 let envelope = Envelope::decode(frame.clone()).unwrap();
1693 let back = Envelope::from_response(envelope.to_response().unwrap()).unwrap();
1694 assert_eq!(back.encode(), frame);
1695 }
1696 }
1697}