1use std::time::{SystemTime, UNIX_EPOCH};
26
27use crate::cbor::{self, Value};
28use crate::identity::KeyPair;
29
30pub const SIG_DOMAIN: &[u8] = b"macula-v2-frame\0";
36
37pub const PROTOCOL_VERSION: i128 = 2;
38
39pub const MAX_FRAME_BYTES: usize = 0x00FF_FFFF;
42
43fn current_millis() -> u64 {
44 SystemTime::now()
45 .duration_since(UNIX_EPOCH)
46 .expect("system clock is after the Unix epoch")
47 .as_millis() as u64
48}
49
50fn fresh_frame_id() -> [u8; 16] {
51 *uuid::Uuid::now_v7().as_bytes()
52}
53
54fn base(
58 frame_type: &str,
59 capabilities: u64,
60 frame_id: [u8; 16],
61 sent_at_ms: u64,
62) -> Vec<(Value, Value)> {
63 vec![
64 (Value::text("version"), Value::Int(PROTOCOL_VERSION)),
65 (Value::text("frame_type"), Value::text(frame_type)),
66 (Value::text("frame_id"), Value::Bytes(frame_id.to_vec())),
67 (Value::text("sent_at_ms"), Value::Int(sent_at_ms as i128)),
68 (
69 Value::text("capabilities"),
70 Value::Int(capabilities as i128),
71 ),
72 (Value::text("realm"), Value::Null),
73 (Value::text("call_id"), Value::Null),
74 (Value::text("source_route"), Value::Null),
75 ]
76}
77
78fn bytes32_list(items: &[[u8; 32]]) -> Value {
79 Value::List(items.iter().map(|b| Value::Bytes(b.to_vec())).collect())
80}
81
82#[derive(Debug, Clone)]
88pub struct ConnectSpec {
89 pub node_id: [u8; 32],
90 pub station_id: [u8; 32],
91 pub realms: Vec<[u8; 32]>,
92 pub capabilities: u64,
93 pub puzzle_evidence: [u8; 32],
94 pub addresses: Vec<Value>,
95 pub site: Option<Value>,
96 pub endorsements: Vec<Value>,
97}
98
99impl ConnectSpec {
100 pub fn new(node_id: [u8; 32], puzzle_evidence: [u8; 32]) -> Self {
105 Self {
106 node_id,
107 station_id: node_id,
110 realms: Vec::new(),
111 capabilities: 0,
112 puzzle_evidence,
113 addresses: Vec::new(),
114 site: None,
115 endorsements: Vec::new(),
116 }
117 }
118}
119
120fn connect_value(spec: &ConnectSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
121 let mut fields = base("connect", spec.capabilities, frame_id, sent_at_ms);
122 fields.push((Value::text("node_id"), Value::Bytes(spec.node_id.to_vec())));
123 fields.push((
124 Value::text("station_id"),
125 Value::Bytes(spec.station_id.to_vec()),
126 ));
127 fields.push((Value::text("realms"), bytes32_list(&spec.realms)));
128 fields.push((
129 Value::text("addresses"),
130 Value::List(spec.addresses.clone()),
131 ));
132 fields.push((
133 Value::text("site"),
134 spec.site.clone().unwrap_or(Value::Null),
135 ));
136 fields.push((
137 Value::text("puzzle_evidence"),
138 Value::Bytes(spec.puzzle_evidence.to_vec()),
139 ));
140 fields.push((
141 Value::text("endorsements"),
142 Value::List(spec.endorsements.clone()),
143 ));
144 Value::Map(fields)
145}
146
147pub fn connect(spec: &ConnectSpec) -> Value {
150 connect_value(spec, fresh_frame_id(), current_millis())
151}
152
153fn goodbye_value(reason: &str, detail: Option<&str>, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
158 let mut fields = base("goodbye", 0, frame_id, sent_at_ms);
159 fields.push((Value::text("reason"), Value::text(reason)));
166 fields.push((
167 Value::text("detail"),
168 detail
169 .map(|d| Value::Bytes(d.as_bytes().to_vec()))
170 .unwrap_or(Value::Null),
171 ));
172 Value::Map(fields)
173}
174
175pub fn goodbye(reason: &str, detail: Option<&str>) -> Value {
178 goodbye_value(reason, detail, fresh_frame_id(), current_millis())
179}
180
181#[derive(Debug, Clone)]
199pub struct CallSpec {
200 pub call_id: [u8; 16],
201 pub procedure: String,
202 pub realm: [u8; 32],
203 pub payload: Value,
204 pub deadline_ms: i128,
205 pub caller: [u8; 32],
206 pub source_route: Vec<u8>,
210 pub retry_budget: u64,
211 pub ucan_token: Vec<u8>,
212}
213
214impl CallSpec {
215 pub fn new(
216 call_id: [u8; 16],
217 procedure: impl Into<String>,
218 realm: [u8; 32],
219 payload: Value,
220 deadline_ms: i128,
221 caller: [u8; 32],
222 ) -> Self {
223 Self {
224 call_id,
225 procedure: procedure.into(),
226 realm,
227 payload,
228 deadline_ms,
229 caller,
230 source_route: Vec::new(),
231 retry_budget: 0,
232 ucan_token: Vec::new(),
233 }
234 }
235}
236
237fn call_value(spec: &CallSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
238 Value::Map(base("call", 0, frame_id, sent_at_ms))
239 .with_field("realm", Value::Bytes(spec.realm.to_vec()))
240 .with_field("call_id", Value::Bytes(spec.call_id.to_vec()))
241 .with_field(
246 "procedure",
247 Value::Bytes(spec.procedure.as_bytes().to_vec()),
248 )
249 .with_field("payload", spec.payload.clone())
250 .with_field("deadline_ms", Value::Int(spec.deadline_ms))
251 .with_field("caller", Value::Bytes(spec.caller.to_vec()))
252 .with_field("source_route", Value::Bytes(spec.source_route.clone()))
253 .with_field("retry_budget", Value::Int(spec.retry_budget as i128))
254 .with_field("ucan_token", Value::Bytes(spec.ucan_token.clone()))
255}
256
257pub fn call(spec: &CallSpec) -> Value {
260 call_value(spec, fresh_frame_id(), current_millis())
261}
262
263#[derive(Debug, Clone)]
265pub struct ResultSpec {
266 pub call_id: [u8; 16],
267 pub payload: Value,
268 pub responded_by: [u8; 32],
269 pub source_route_reverse: Vec<u8>,
270}
271
272impl ResultSpec {
273 pub fn new(call_id: [u8; 16], payload: Value, responded_by: [u8; 32]) -> Self {
274 Self {
275 call_id,
276 payload,
277 responded_by,
278 source_route_reverse: Vec::new(),
279 }
280 }
281}
282
283fn result_value(spec: &ResultSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
284 Value::Map(base("result", 0, frame_id, sent_at_ms))
290 .with_field("call_id", Value::Bytes(spec.call_id.to_vec()))
291 .with_field("payload", spec.payload.clone())
292 .with_field("responded_by", Value::Bytes(spec.responded_by.to_vec()))
293 .with_field(
294 "source_route_reverse",
295 Value::Bytes(spec.source_route_reverse.clone()),
296 )
297}
298
299pub fn result(spec: &ResultSpec) -> Value {
301 result_value(spec, fresh_frame_id(), current_millis())
302}
303
304#[derive(Debug, Clone)]
308pub struct CallErrorSpec {
309 pub call_id: [u8; 16],
310 pub code: crate::bolt4::Code,
311 pub reported_by: [u8; 32],
312 pub detail: Option<String>,
313 pub offending_hop: Option<[u8; 32]>,
314 pub source_route_partial: Vec<u8>,
315}
316
317impl CallErrorSpec {
318 pub fn new(call_id: [u8; 16], code: crate::bolt4::Code, reported_by: [u8; 32]) -> Self {
319 Self {
320 call_id,
321 code,
322 reported_by,
323 detail: None,
324 offending_hop: None,
325 source_route_partial: Vec::new(),
326 }
327 }
328}
329
330fn call_error_value(spec: &CallErrorSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
331 Value::Map(base("error", 0, frame_id, sent_at_ms))
332 .with_field("call_id", Value::Bytes(spec.call_id.to_vec()))
333 .with_field("code", Value::Int(spec.code.as_u8() as i128))
334 .with_field("name", Value::text(spec.code.name()))
335 .with_field("reported_by", Value::Bytes(spec.reported_by.to_vec()))
336 .with_field(
337 "detail",
340 spec.detail
341 .as_ref()
342 .map(|d| Value::Bytes(d.as_bytes().to_vec()))
343 .unwrap_or(Value::Null),
344 )
345 .with_field(
346 "offending_hop",
347 spec.offending_hop
348 .map(|h| Value::Bytes(h.to_vec()))
349 .unwrap_or(Value::Null),
350 )
351 .with_field(
352 "source_route_partial",
353 Value::Bytes(spec.source_route_partial.clone()),
354 )
355}
356
357pub fn call_error(spec: &CallErrorSpec) -> Value {
359 call_error_value(spec, fresh_frame_id(), current_millis())
360}
361
362#[derive(Debug, Clone)]
369pub struct CallInfo {
370 pub call_id: [u8; 16],
371 pub procedure: String,
372 pub realm: [u8; 32],
373 pub payload: Value,
374 pub deadline_ms: i128,
375 pub caller: [u8; 32],
376 pub ucan_token: Vec<u8>,
379}
380
381#[derive(Debug, PartialEq, Eq)]
382pub enum ParseCallError {
383 NotACallFrame,
384 MissingField(&'static str),
385 WrongFieldType(&'static str),
386}
387
388impl std::fmt::Display for ParseCallError {
389 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
390 match self {
391 ParseCallError::NotACallFrame => write!(f, "frame_type is not \"call\""),
392 ParseCallError::MissingField(name) => write!(f, "missing required field {name:?}"),
393 ParseCallError::WrongFieldType(name) => write!(f, "field {name:?} has the wrong type"),
394 }
395 }
396}
397
398impl std::error::Error for ParseCallError {}
399
400pub fn parse_call(frame: &Value) -> Result<CallInfo, ParseCallError> {
403 match frame.get("frame_type") {
404 Some(Value::Text(t)) if t == "call" => {}
405 _ => return Err(ParseCallError::NotACallFrame),
406 }
407 let call_id = match frame.get("call_id") {
408 Some(Value::Bytes(b)) => b
409 .as_slice()
410 .try_into()
411 .map_err(|_| ParseCallError::WrongFieldType("call_id"))?,
412 Some(_) => return Err(ParseCallError::WrongFieldType("call_id")),
413 None => return Err(ParseCallError::MissingField("call_id")),
414 };
415 let procedure = match frame.get("procedure") {
417 Some(Value::Bytes(b)) => {
418 String::from_utf8(b.clone()).map_err(|_| ParseCallError::WrongFieldType("procedure"))?
419 }
420 Some(_) => return Err(ParseCallError::WrongFieldType("procedure")),
421 None => return Err(ParseCallError::MissingField("procedure")),
422 };
423 let realm = match frame.get("realm") {
424 Some(Value::Bytes(b)) => b
425 .as_slice()
426 .try_into()
427 .map_err(|_| ParseCallError::WrongFieldType("realm"))?,
428 Some(_) => return Err(ParseCallError::WrongFieldType("realm")),
429 None => return Err(ParseCallError::MissingField("realm")),
430 };
431 let payload = frame
432 .get("payload")
433 .cloned()
434 .ok_or(ParseCallError::MissingField("payload"))?;
435 let deadline_ms = match frame.get("deadline_ms") {
436 Some(Value::Int(n)) => *n,
437 Some(_) => return Err(ParseCallError::WrongFieldType("deadline_ms")),
438 None => return Err(ParseCallError::MissingField("deadline_ms")),
439 };
440 let caller = match frame.get("caller") {
441 Some(Value::Bytes(b)) => b
442 .as_slice()
443 .try_into()
444 .map_err(|_| ParseCallError::WrongFieldType("caller"))?,
445 Some(_) => return Err(ParseCallError::WrongFieldType("caller")),
446 None => return Err(ParseCallError::MissingField("caller")),
447 };
448 let ucan_token = match frame.get("ucan_token") {
449 Some(Value::Bytes(b)) => b.clone(),
450 _ => Vec::new(),
451 };
452 Ok(CallInfo {
453 call_id,
454 procedure,
455 realm,
456 payload,
457 deadline_ms,
458 caller,
459 ucan_token,
460 })
461}
462
463#[derive(Debug, Clone)]
466pub enum CallResponse {
467 Result {
468 payload: Value,
469 responded_by: [u8; 32],
470 },
471 Error {
472 code: u8,
473 name: String,
474 reported_by: [u8; 32],
475 detail: Option<String>,
476 },
477}
478
479#[derive(Debug, PartialEq, Eq)]
480pub enum ParseCallResponseError {
481 NotAResultOrError,
482 MissingField(&'static str),
483 WrongFieldType(&'static str),
484}
485
486impl std::fmt::Display for ParseCallResponseError {
487 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
488 match self {
489 ParseCallResponseError::NotAResultOrError => {
490 write!(f, "frame_type is neither \"result\" nor \"error\"")
491 }
492 ParseCallResponseError::MissingField(name) => {
493 write!(f, "missing required field {name:?}")
494 }
495 ParseCallResponseError::WrongFieldType(name) => {
496 write!(f, "field {name:?} has the wrong type")
497 }
498 }
499 }
500}
501
502impl std::error::Error for ParseCallResponseError {}
503
504pub fn frame_call_id(frame: &Value) -> Option<[u8; 16]> {
510 match frame.get("call_id") {
511 Some(Value::Bytes(b)) => b.as_slice().try_into().ok(),
512 _ => None,
513 }
514}
515
516pub fn parse_call_response(frame: &Value) -> Result<CallResponse, ParseCallResponseError> {
518 match frame.get("frame_type") {
519 Some(Value::Text(t)) if t == "result" => {
520 let payload = frame
521 .get("payload")
522 .cloned()
523 .ok_or(ParseCallResponseError::MissingField("payload"))?;
524 let responded_by = get_bytes32_generic(frame, "responded_by")?;
525 Ok(CallResponse::Result {
526 payload,
527 responded_by,
528 })
529 }
530 Some(Value::Text(t)) if t == "error" => {
531 let code = match frame.get("code") {
532 Some(Value::Int(n)) if (0..=255).contains(n) => *n as u8,
533 Some(_) => return Err(ParseCallResponseError::WrongFieldType("code")),
534 None => return Err(ParseCallResponseError::MissingField("code")),
535 };
536 let name = match frame.get("name") {
537 Some(Value::Text(t)) => t.clone(),
538 Some(_) => return Err(ParseCallResponseError::WrongFieldType("name")),
539 None => return Err(ParseCallResponseError::MissingField("name")),
540 };
541 let reported_by = get_bytes32_generic(frame, "reported_by")?;
542 let detail = match frame.get("detail") {
545 None | Some(Value::Null) => None,
546 Some(Value::Bytes(b)) => Some(
547 String::from_utf8(b.clone())
548 .map_err(|_| ParseCallResponseError::WrongFieldType("detail"))?,
549 ),
550 Some(_) => return Err(ParseCallResponseError::WrongFieldType("detail")),
551 };
552 Ok(CallResponse::Error {
553 code,
554 name,
555 reported_by,
556 detail,
557 })
558 }
559 _ => Err(ParseCallResponseError::NotAResultOrError),
560 }
561}
562
563fn get_bytes32_generic(
564 frame: &Value,
565 field: &'static str,
566) -> Result<[u8; 32], ParseCallResponseError> {
567 match frame.get(field) {
568 None => Err(ParseCallResponseError::MissingField(field)),
569 Some(Value::Bytes(b)) => b
570 .as_slice()
571 .try_into()
572 .map_err(|_| ParseCallResponseError::WrongFieldType(field)),
573 Some(_) => Err(ParseCallResponseError::WrongFieldType(field)),
574 }
575}
576
577#[derive(Debug, Clone)]
583pub struct PublishSpec {
584 pub topic: String,
585 pub realm: [u8; 32],
586 pub publisher: [u8; 32],
587 pub seq: u64,
588 pub payload: Value,
589 pub published_at_ms: u64,
590 pub ttl_ms: Option<u64>,
591}
592
593impl PublishSpec {
594 pub fn new(
595 topic: impl Into<String>,
596 realm: [u8; 32],
597 publisher: [u8; 32],
598 seq: u64,
599 payload: Value,
600 published_at_ms: u64,
601 ) -> Self {
602 Self {
603 topic: topic.into(),
604 realm,
605 publisher,
606 seq,
607 payload,
608 published_at_ms,
609 ttl_ms: None,
610 }
611 }
612}
613
614fn publish_value(spec: &PublishSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
615 Value::Map(base("publish", 0, frame_id, sent_at_ms))
616 .with_field("realm", Value::Bytes(spec.realm.to_vec()))
617 .with_field("topic", Value::Bytes(spec.topic.as_bytes().to_vec()))
620 .with_field("publisher", Value::Bytes(spec.publisher.to_vec()))
621 .with_field("seq", Value::Int(spec.seq as i128))
622 .with_field("payload", spec.payload.clone())
623 .with_field("published_at_ms", Value::Int(spec.published_at_ms as i128))
624 .with_field(
625 "ttl_ms",
626 spec.ttl_ms
627 .map(|t| Value::Int(t as i128))
628 .unwrap_or(Value::Null),
629 )
630}
631
632pub fn publish(spec: &PublishSpec) -> Value {
636 publish_value(spec, fresh_frame_id(), current_millis())
637}
638
639#[derive(Debug, Clone)]
641pub struct SubscribeSpec {
642 pub topic: String,
643 pub realm: [u8; 32],
644 pub subscriber: [u8; 32],
645}
646
647impl SubscribeSpec {
648 pub fn new(topic: impl Into<String>, realm: [u8; 32], subscriber: [u8; 32]) -> Self {
649 Self {
650 topic: topic.into(),
651 realm,
652 subscriber,
653 }
654 }
655}
656
657fn subscribe_value(spec: &SubscribeSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
658 Value::Map(base("subscribe", 0, frame_id, sent_at_ms))
659 .with_field("realm", Value::Bytes(spec.realm.to_vec()))
660 .with_field("topic", Value::Bytes(spec.topic.as_bytes().to_vec()))
663 .with_field("subscriber", Value::Bytes(spec.subscriber.to_vec()))
664 .with_field("filter", Value::Null)
665 .with_field("options", Value::Map(vec![]))
666}
667
668pub fn subscribe(spec: &SubscribeSpec) -> Value {
671 subscribe_value(spec, fresh_frame_id(), current_millis())
672}
673
674#[derive(Debug, Clone)]
676pub struct UnsubscribeSpec {
677 pub topic: String,
678 pub realm: [u8; 32],
679 pub subscriber: [u8; 32],
680}
681
682impl UnsubscribeSpec {
683 pub fn new(topic: impl Into<String>, realm: [u8; 32], subscriber: [u8; 32]) -> Self {
684 Self {
685 topic: topic.into(),
686 realm,
687 subscriber,
688 }
689 }
690}
691
692fn unsubscribe_value(spec: &UnsubscribeSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
693 Value::Map(base("unsubscribe", 0, frame_id, sent_at_ms))
694 .with_field("realm", Value::Bytes(spec.realm.to_vec()))
695 .with_field("topic", Value::Bytes(spec.topic.as_bytes().to_vec()))
698 .with_field("subscriber", Value::Bytes(spec.subscriber.to_vec()))
699}
700
701pub fn unsubscribe(spec: &UnsubscribeSpec) -> Value {
703 unsubscribe_value(spec, fresh_frame_id(), current_millis())
704}
705
706#[derive(Debug, Clone)]
708pub struct EventInfo {
709 pub topic: String,
710 pub realm: [u8; 32],
711 pub publisher: [u8; 32],
712 pub seq: u64,
713 pub payload: Value,
714 pub delivered_via: String,
715}
716
717#[derive(Debug, PartialEq, Eq)]
718pub enum ParseEventError {
719 NotAnEventFrame,
720 MissingField(&'static str),
721 WrongFieldType(&'static str),
722}
723
724impl std::fmt::Display for ParseEventError {
725 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
726 match self {
727 ParseEventError::NotAnEventFrame => write!(f, "frame_type is not \"event\""),
728 ParseEventError::MissingField(name) => write!(f, "missing required field {name:?}"),
729 ParseEventError::WrongFieldType(name) => write!(f, "field {name:?} has the wrong type"),
730 }
731 }
732}
733
734impl std::error::Error for ParseEventError {}
735
736pub fn parse_event(frame: &Value) -> Result<EventInfo, ParseEventError> {
738 match frame.get("frame_type") {
739 Some(Value::Text(t)) if t == "event" => {}
740 _ => return Err(ParseEventError::NotAnEventFrame),
741 }
742 let topic = match frame.get("topic") {
744 Some(Value::Bytes(b)) => {
745 String::from_utf8(b.clone()).map_err(|_| ParseEventError::WrongFieldType("topic"))?
746 }
747 Some(_) => return Err(ParseEventError::WrongFieldType("topic")),
748 None => return Err(ParseEventError::MissingField("topic")),
749 };
750 let realm = match frame.get("realm") {
751 Some(Value::Bytes(b)) => b
752 .as_slice()
753 .try_into()
754 .map_err(|_| ParseEventError::WrongFieldType("realm"))?,
755 Some(_) => return Err(ParseEventError::WrongFieldType("realm")),
756 None => return Err(ParseEventError::MissingField("realm")),
757 };
758 let publisher = match frame.get("publisher") {
759 Some(Value::Bytes(b)) => b
760 .as_slice()
761 .try_into()
762 .map_err(|_| ParseEventError::WrongFieldType("publisher"))?,
763 Some(_) => return Err(ParseEventError::WrongFieldType("publisher")),
764 None => return Err(ParseEventError::MissingField("publisher")),
765 };
766 let seq = match frame.get("seq") {
767 Some(Value::Int(n)) if *n >= 0 => *n as u64,
768 Some(_) => return Err(ParseEventError::WrongFieldType("seq")),
769 None => return Err(ParseEventError::MissingField("seq")),
770 };
771 let payload = frame
772 .get("payload")
773 .cloned()
774 .ok_or(ParseEventError::MissingField("payload"))?;
775 let delivered_via = match frame.get("delivered_via") {
776 Some(Value::Text(t)) => t.clone(),
777 Some(_) => return Err(ParseEventError::WrongFieldType("delivered_via")),
778 None => return Err(ParseEventError::MissingField("delivered_via")),
779 };
780 Ok(EventInfo {
781 topic,
782 realm,
783 publisher,
784 seq,
785 payload,
786 delivered_via,
787 })
788}
789
790#[derive(Debug, Clone, PartialEq, Eq)]
797pub struct HelloInfo {
798 pub node_id: [u8; 32],
799 pub station_id: [u8; 32],
800 pub realms: Vec<[u8; 32]>,
801 pub capabilities: u64,
802 pub accepted: bool,
803 pub negotiated_capabilities: u64,
804 pub refusal_code: Option<i128>,
805}
806
807#[derive(Debug, PartialEq, Eq)]
808pub enum ParseHelloError {
809 NotAHelloFrame,
810 MissingField(&'static str),
811 WrongFieldType(&'static str),
812}
813
814impl std::fmt::Display for ParseHelloError {
815 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
816 match self {
817 ParseHelloError::NotAHelloFrame => write!(f, "frame_type is not \"hello\""),
818 ParseHelloError::MissingField(name) => write!(f, "missing required field {name:?}"),
819 ParseHelloError::WrongFieldType(name) => write!(f, "field {name:?} has the wrong type"),
820 }
821 }
822}
823
824impl std::error::Error for ParseHelloError {}
825
826fn get_bytes32(frame: &Value, field: &'static str) -> Result<[u8; 32], ParseHelloError> {
827 match frame.get(field) {
828 None => Err(ParseHelloError::MissingField(field)),
829 Some(Value::Bytes(b)) => b
830 .as_slice()
831 .try_into()
832 .map_err(|_| ParseHelloError::WrongFieldType(field)),
833 Some(_) => Err(ParseHelloError::WrongFieldType(field)),
834 }
835}
836
837fn get_bytes32_list(frame: &Value, field: &'static str) -> Result<Vec<[u8; 32]>, ParseHelloError> {
838 match frame.get(field) {
839 None => Err(ParseHelloError::MissingField(field)),
840 Some(Value::List(items)) => items
841 .iter()
842 .map(|v| match v {
843 Value::Bytes(b) => b
844 .as_slice()
845 .try_into()
846 .map_err(|_| ParseHelloError::WrongFieldType(field)),
847 _ => Err(ParseHelloError::WrongFieldType(field)),
848 })
849 .collect(),
850 Some(_) => Err(ParseHelloError::WrongFieldType(field)),
851 }
852}
853
854fn get_uint(frame: &Value, field: &'static str) -> Result<u64, ParseHelloError> {
855 match frame.get(field) {
856 None => Err(ParseHelloError::MissingField(field)),
857 Some(Value::Int(n)) if *n >= 0 => Ok(*n as u64),
858 Some(_) => Err(ParseHelloError::WrongFieldType(field)),
859 }
860}
861
862fn get_bool(frame: &Value, field: &'static str) -> Result<bool, ParseHelloError> {
863 match frame.get(field) {
864 None => Err(ParseHelloError::MissingField(field)),
865 Some(Value::Text(t)) if t == "true" => Ok(true),
866 Some(Value::Text(t)) if t == "false" => Ok(false),
867 Some(_) => Err(ParseHelloError::WrongFieldType(field)),
868 }
869}
870
871pub fn parse_hello(frame: &Value) -> Result<HelloInfo, ParseHelloError> {
873 match frame.get("frame_type") {
874 Some(Value::Text(t)) if t == "hello" => {}
875 _ => return Err(ParseHelloError::NotAHelloFrame),
876 }
877 let refusal_code = match frame.get("refusal_code") {
878 None | Some(Value::Null) => None,
879 Some(Value::Int(n)) => Some(*n),
880 Some(_) => return Err(ParseHelloError::WrongFieldType("refusal_code")),
881 };
882 Ok(HelloInfo {
883 node_id: get_bytes32(frame, "node_id")?,
884 station_id: get_bytes32(frame, "station_id")?,
885 realms: get_bytes32_list(frame, "realms")?,
886 capabilities: get_uint(frame, "capabilities")?,
887 accepted: get_bool(frame, "accepted")?,
888 negotiated_capabilities: get_uint(frame, "negotiated_capabilities")?,
889 refusal_code,
890 })
891}
892
893pub fn sign(frame: Value, identity: &KeyPair) -> Value {
901 let signable = signable_bytes(&frame);
902 let sig = identity.sign(&signable);
903 frame.with_field("signature", Value::Bytes(sig.to_vec()))
904}
905
906fn signable_bytes(frame: &Value) -> Vec<u8> {
907 let unsigned = frame.without(&["signature", "publisher_sig"]);
908 let canonical =
909 cbor::encode(&unsigned).expect("a frame built by this module is always encodable");
910 let mut out = Vec::with_capacity(SIG_DOMAIN.len() + canonical.len());
911 out.extend_from_slice(SIG_DOMAIN);
912 out.extend_from_slice(&canonical);
913 out
914}
915
916#[derive(Debug, PartialEq, Eq)]
917pub enum VerifyError {
918 MissingSignature,
919 BadSignature,
920 SignatureInvalid,
921}
922
923impl std::fmt::Display for VerifyError {
924 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
925 match self {
926 VerifyError::MissingSignature => write!(f, "frame has no signature field"),
927 VerifyError::BadSignature => write!(f, "signature field is not 64 bytes"),
928 VerifyError::SignatureInvalid => write!(f, "signature does not verify against pubkey"),
929 }
930 }
931}
932
933impl std::error::Error for VerifyError {}
934
935pub fn verify(frame: &Value, pubkey: &[u8; 32]) -> Result<(), VerifyError> {
938 let sig: [u8; 64] = match frame.get("signature") {
939 Some(Value::Bytes(b)) => b
940 .as_slice()
941 .try_into()
942 .map_err(|_| VerifyError::BadSignature)?,
943 _ => return Err(VerifyError::MissingSignature),
944 };
945 let signable = signable_bytes(frame);
946 if crate::identity::verify(&signable, &sig, pubkey) {
947 Ok(())
948 } else {
949 Err(VerifyError::SignatureInvalid)
950 }
951}
952
953pub const EVENT_PUBLISHER_DOMAIN: &[u8] = b"macula-v2-event-pub\0";
971
972pub fn sign_publisher(frame: Value, identity: &KeyPair) -> Value {
978 let signable = publisher_signing_bytes(&frame);
979 let sig = identity.sign(&signable);
980 frame.with_field("publisher_sig", Value::Bytes(sig.to_vec()))
981}
982
983#[derive(Debug, PartialEq, Eq)]
984pub enum VerifyPublisherError {
985 MissingPublisherSig,
986 BadPublisherSig,
987 PublisherSigInvalid,
988}
989
990impl std::fmt::Display for VerifyPublisherError {
991 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
992 match self {
993 VerifyPublisherError::MissingPublisherSig => {
994 write!(f, "frame has no publisher_sig field")
995 }
996 VerifyPublisherError::BadPublisherSig => {
997 write!(f, "publisher_sig field is not 64 bytes")
998 }
999 VerifyPublisherError::PublisherSigInvalid => write!(
1000 f,
1001 "publisher_sig does not verify against the frame's publisher field"
1002 ),
1003 }
1004 }
1005}
1006
1007impl std::error::Error for VerifyPublisherError {}
1008
1009pub fn verify_publisher(frame: &Value) -> Result<(), VerifyPublisherError> {
1015 let sig: [u8; 64] = match frame.get("publisher_sig") {
1016 Some(Value::Bytes(b)) => b
1017 .as_slice()
1018 .try_into()
1019 .map_err(|_| VerifyPublisherError::BadPublisherSig)?,
1020 _ => return Err(VerifyPublisherError::MissingPublisherSig),
1021 };
1022 let pubkey: [u8; 32] = match frame.get("publisher") {
1023 Some(Value::Bytes(b)) => b
1024 .as_slice()
1025 .try_into()
1026 .map_err(|_| VerifyPublisherError::BadPublisherSig)?,
1027 _ => return Err(VerifyPublisherError::BadPublisherSig),
1028 };
1029 let signable = publisher_signing_bytes(frame);
1030 if crate::identity::verify(&signable, &sig, &pubkey) {
1031 Ok(())
1032 } else {
1033 Err(VerifyPublisherError::PublisherSigInvalid)
1034 }
1035}
1036
1037fn publisher_signing_bytes(frame: &Value) -> Vec<u8> {
1042 let fields = ["topic", "realm", "publisher", "seq", "payload"];
1043 let pairs: Vec<(Value, Value)> = fields
1044 .iter()
1045 .map(|f| {
1046 let v = frame.get(f).cloned().unwrap_or(Value::Null);
1047 (Value::text(*f), v)
1048 })
1049 .collect();
1050 let canonical =
1051 cbor::encode(&Value::Map(pairs)).expect("a frame built by this module is always encodable");
1052 let mut out = Vec::with_capacity(EVENT_PUBLISHER_DOMAIN.len() + canonical.len());
1053 out.extend_from_slice(EVENT_PUBLISHER_DOMAIN);
1054 out.extend_from_slice(&canonical);
1055 out
1056}
1057
1058#[derive(Debug, PartialEq, Eq)]
1063pub enum EncodeFrameError {
1064 TooLarge(usize),
1065 Cbor(cbor::IntOutOfRange),
1066}
1067
1068impl std::fmt::Display for EncodeFrameError {
1069 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1070 match self {
1071 EncodeFrameError::TooLarge(n) => {
1072 write!(
1073 f,
1074 "frame is {n} bytes, exceeding the {MAX_FRAME_BYTES}-byte cap"
1075 )
1076 }
1077 EncodeFrameError::Cbor(e) => write!(f, "{e}"),
1078 }
1079 }
1080}
1081
1082impl std::error::Error for EncodeFrameError {}
1083
1084pub fn encode(frame: &Value) -> Result<Vec<u8>, EncodeFrameError> {
1086 let payload = cbor::encode(frame).map_err(EncodeFrameError::Cbor)?;
1087 if payload.len() > MAX_FRAME_BYTES {
1088 return Err(EncodeFrameError::TooLarge(payload.len()));
1089 }
1090 let mut out = Vec::with_capacity(4 + payload.len());
1091 out.extend_from_slice(&(payload.len() as u32).to_be_bytes());
1092 out.extend_from_slice(&payload);
1093 Ok(out)
1094}
1095
1096#[derive(Debug)]
1101pub enum Decoded {
1102 Frame(Value, usize),
1105 More(usize),
1108}
1109
1110#[derive(Debug)]
1111pub enum DecodeFrameError {
1112 TooLarge(usize),
1113 Cbor(cbor::DecodeError),
1114}
1115
1116impl std::fmt::Display for DecodeFrameError {
1117 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1118 match self {
1119 DecodeFrameError::TooLarge(n) => {
1120 write!(
1121 f,
1122 "claimed frame length {n} exceeds the {MAX_FRAME_BYTES}-byte cap"
1123 )
1124 }
1125 DecodeFrameError::Cbor(e) => write!(f, "{e}"),
1126 }
1127 }
1128}
1129
1130impl std::error::Error for DecodeFrameError {}
1131
1132pub fn decode(buf: &[u8]) -> Result<Decoded, DecodeFrameError> {
1134 if buf.len() < 4 {
1135 return Ok(Decoded::More(4 - buf.len()));
1136 }
1137 let len = u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]) as usize;
1138 if len > MAX_FRAME_BYTES {
1139 return Err(DecodeFrameError::TooLarge(len));
1140 }
1141 if buf.len() < 4 + len {
1142 return Ok(Decoded::More(4 + len - buf.len()));
1143 }
1144 let value = cbor::decode(&buf[4..4 + len]).map_err(DecodeFrameError::Cbor)?;
1145 Ok(Decoded::Frame(value, 4 + len))
1146}
1147
1148#[derive(Debug, Clone)]
1159pub struct AdvertiseSpec {
1160 pub realm: [u8; 32],
1161 pub procedure: String,
1162 pub advertiser: [u8; 32],
1163}
1164
1165impl AdvertiseSpec {
1166 pub fn new(realm: [u8; 32], procedure: impl Into<String>, advertiser: [u8; 32]) -> Self {
1167 Self {
1168 realm,
1169 procedure: procedure.into(),
1170 advertiser,
1171 }
1172 }
1173}
1174
1175fn advertise_value(spec: &AdvertiseSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
1176 Value::Map(base("advertise", 0, frame_id, sent_at_ms))
1181 .with_field("realm", Value::Bytes(spec.realm.to_vec()))
1182 .with_field(
1185 "procedure",
1186 Value::Bytes(spec.procedure.as_bytes().to_vec()),
1187 )
1188 .with_field("advertiser", Value::Bytes(spec.advertiser.to_vec()))
1189 .with_field("options", Value::Map(vec![]))
1192}
1193
1194pub fn advertise(spec: &AdvertiseSpec) -> Value {
1196 advertise_value(spec, fresh_frame_id(), current_millis())
1197}
1198
1199#[derive(Debug, Clone)]
1201pub struct UnadvertiseSpec {
1202 pub realm: [u8; 32],
1203 pub procedure: String,
1204 pub advertiser: [u8; 32],
1205}
1206
1207impl UnadvertiseSpec {
1208 pub fn new(realm: [u8; 32], procedure: impl Into<String>, advertiser: [u8; 32]) -> Self {
1209 Self {
1210 realm,
1211 procedure: procedure.into(),
1212 advertiser,
1213 }
1214 }
1215}
1216
1217fn unadvertise_value(spec: &UnadvertiseSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
1218 Value::Map(base("unadvertise", 0, frame_id, sent_at_ms))
1219 .with_field("realm", Value::Bytes(spec.realm.to_vec()))
1220 .with_field(
1221 "procedure",
1222 Value::Bytes(spec.procedure.as_bytes().to_vec()),
1223 )
1224 .with_field("advertiser", Value::Bytes(spec.advertiser.to_vec()))
1225}
1226
1227pub fn unadvertise(spec: &UnadvertiseSpec) -> Value {
1229 unadvertise_value(spec, fresh_frame_id(), current_millis())
1230}
1231
1232#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1291pub enum StreamMode {
1292 ServerStream,
1294 ClientStream,
1297 Bidi,
1299}
1300
1301impl StreamMode {
1302 pub fn name(self) -> &'static str {
1303 match self {
1304 StreamMode::ServerStream => "server_stream",
1305 StreamMode::ClientStream => "client_stream",
1306 StreamMode::Bidi => "bidi",
1307 }
1308 }
1309
1310 fn from_name(name: &str) -> Option<Self> {
1311 match name {
1312 "server_stream" => Some(StreamMode::ServerStream),
1313 "client_stream" => Some(StreamMode::ClientStream),
1314 "bidi" => Some(StreamMode::Bidi),
1315 _ => None,
1316 }
1317 }
1318}
1319
1320#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1323pub enum StreamEncoding {
1324 Raw,
1326 Msgpack,
1329}
1330
1331impl StreamEncoding {
1332 pub fn name(self) -> &'static str {
1333 match self {
1334 StreamEncoding::Raw => "raw",
1335 StreamEncoding::Msgpack => "msgpack",
1336 }
1337 }
1338
1339 fn from_name(name: &str) -> Option<Self> {
1340 match name {
1341 "raw" => Some(StreamEncoding::Raw),
1342 "msgpack" => Some(StreamEncoding::Msgpack),
1343 _ => None,
1344 }
1345 }
1346}
1347
1348#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1350pub enum StreamRole {
1351 Send,
1353 Both,
1355}
1356
1357impl StreamRole {
1358 pub fn name(self) -> &'static str {
1359 match self {
1360 StreamRole::Send => "send",
1361 StreamRole::Both => "both",
1362 }
1363 }
1364
1365 fn from_name(name: &str) -> Option<Self> {
1366 match name {
1367 "send" => Some(StreamRole::Send),
1368 "both" => Some(StreamRole::Both),
1369 _ => None,
1370 }
1371 }
1372}
1373
1374#[derive(Debug, Clone)]
1378pub struct StreamOpenSpec {
1379 pub stream_id: [u8; 16],
1380 pub procedure: String,
1381 pub realm: [u8; 32],
1382 pub mode: StreamMode,
1383 pub args: Value,
1384 pub deadline_ms: i128,
1385 pub caller: [u8; 32],
1386 pub source_route: Vec<u8>,
1387 pub retry_budget: u64,
1388}
1389
1390impl StreamOpenSpec {
1391 pub fn new(
1392 stream_id: [u8; 16],
1393 procedure: impl Into<String>,
1394 realm: [u8; 32],
1395 mode: StreamMode,
1396 args: Value,
1397 deadline_ms: i128,
1398 caller: [u8; 32],
1399 ) -> Self {
1400 Self {
1401 stream_id,
1402 procedure: procedure.into(),
1403 realm,
1404 mode,
1405 args,
1406 deadline_ms,
1407 caller,
1408 source_route: Vec::new(),
1409 retry_budget: 0,
1410 }
1411 }
1412}
1413
1414fn stream_open_value(spec: &StreamOpenSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
1415 Value::Map(base("stream_open", 0, frame_id, sent_at_ms))
1416 .with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
1417 .with_field(
1420 "procedure",
1421 Value::Bytes(spec.procedure.as_bytes().to_vec()),
1422 )
1423 .with_field("realm", Value::Bytes(spec.realm.to_vec()))
1424 .with_field("mode", Value::text(spec.mode.name()))
1425 .with_field("args", spec.args.clone())
1426 .with_field("deadline_ms", Value::Int(spec.deadline_ms))
1427 .with_field("caller", Value::Bytes(spec.caller.to_vec()))
1428 .with_field("source_route", Value::Bytes(spec.source_route.clone()))
1429 .with_field("retry_budget", Value::Int(spec.retry_budget as i128))
1430}
1431
1432pub fn stream_open(spec: &StreamOpenSpec) -> Value {
1435 stream_open_value(spec, fresh_frame_id(), current_millis())
1436}
1437
1438#[derive(Debug, Clone)]
1443pub struct StreamOpenInfo {
1444 pub stream_id: [u8; 16],
1445 pub procedure: String,
1446 pub realm: [u8; 32],
1447 pub mode: StreamMode,
1448 pub args: Value,
1449 pub deadline_ms: i128,
1450 pub caller: [u8; 32],
1451}
1452
1453#[derive(Debug, PartialEq, Eq)]
1454pub enum ParseStreamOpenError {
1455 NotAStreamOpenFrame,
1456 MissingField(&'static str),
1457 WrongFieldType(&'static str),
1458}
1459
1460impl std::fmt::Display for ParseStreamOpenError {
1461 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1462 match self {
1463 ParseStreamOpenError::NotAStreamOpenFrame => {
1464 write!(f, "frame_type is not \"stream_open\"")
1465 }
1466 ParseStreamOpenError::MissingField(name) => {
1467 write!(f, "missing required field {name:?}")
1468 }
1469 ParseStreamOpenError::WrongFieldType(name) => {
1470 write!(f, "field {name:?} has the wrong type")
1471 }
1472 }
1473 }
1474}
1475
1476impl std::error::Error for ParseStreamOpenError {}
1477
1478pub fn parse_stream_open(frame: &Value) -> Result<StreamOpenInfo, ParseStreamOpenError> {
1480 match frame.get("frame_type") {
1481 Some(Value::Text(t)) if t == "stream_open" => {}
1482 _ => return Err(ParseStreamOpenError::NotAStreamOpenFrame),
1483 }
1484 let stream_id = match frame.get("stream_id") {
1485 Some(Value::Bytes(b)) => b
1486 .as_slice()
1487 .try_into()
1488 .map_err(|_| ParseStreamOpenError::WrongFieldType("stream_id"))?,
1489 Some(_) => return Err(ParseStreamOpenError::WrongFieldType("stream_id")),
1490 None => return Err(ParseStreamOpenError::MissingField("stream_id")),
1491 };
1492 let procedure = match frame.get("procedure") {
1494 Some(Value::Bytes(b)) => String::from_utf8(b.clone())
1495 .map_err(|_| ParseStreamOpenError::WrongFieldType("procedure"))?,
1496 Some(_) => return Err(ParseStreamOpenError::WrongFieldType("procedure")),
1497 None => return Err(ParseStreamOpenError::MissingField("procedure")),
1498 };
1499 let realm = match frame.get("realm") {
1500 Some(Value::Bytes(b)) => b
1501 .as_slice()
1502 .try_into()
1503 .map_err(|_| ParseStreamOpenError::WrongFieldType("realm"))?,
1504 Some(_) => return Err(ParseStreamOpenError::WrongFieldType("realm")),
1505 None => return Err(ParseStreamOpenError::MissingField("realm")),
1506 };
1507 let mode = match frame.get("mode") {
1508 Some(Value::Text(t)) => {
1509 StreamMode::from_name(t).ok_or(ParseStreamOpenError::WrongFieldType("mode"))?
1510 }
1511 Some(_) => return Err(ParseStreamOpenError::WrongFieldType("mode")),
1512 None => return Err(ParseStreamOpenError::MissingField("mode")),
1513 };
1514 let args = frame
1515 .get("args")
1516 .cloned()
1517 .ok_or(ParseStreamOpenError::MissingField("args"))?;
1518 let deadline_ms = match frame.get("deadline_ms") {
1519 Some(Value::Int(n)) => *n,
1520 Some(_) => return Err(ParseStreamOpenError::WrongFieldType("deadline_ms")),
1521 None => return Err(ParseStreamOpenError::MissingField("deadline_ms")),
1522 };
1523 let caller = match frame.get("caller") {
1524 Some(Value::Bytes(b)) => b
1525 .as_slice()
1526 .try_into()
1527 .map_err(|_| ParseStreamOpenError::WrongFieldType("caller"))?,
1528 Some(_) => return Err(ParseStreamOpenError::WrongFieldType("caller")),
1529 None => return Err(ParseStreamOpenError::MissingField("caller")),
1530 };
1531 Ok(StreamOpenInfo {
1532 stream_id,
1533 procedure,
1534 realm,
1535 mode,
1536 args,
1537 deadline_ms,
1538 caller,
1539 })
1540}
1541
1542#[derive(Debug, Clone)]
1556pub struct StreamDataSpec {
1557 pub stream_id: [u8; 16],
1558 pub seq: u64,
1559 pub encoding: StreamEncoding,
1560 pub body: Value,
1561 pub signer: Option<[u8; 32]>,
1562}
1563
1564impl StreamDataSpec {
1565 pub fn new(
1566 stream_id: [u8; 16],
1567 seq: u64,
1568 encoding: StreamEncoding,
1569 body: Value,
1570 signer: Option<[u8; 32]>,
1571 ) -> Self {
1572 Self {
1573 stream_id,
1574 seq,
1575 encoding,
1576 body,
1577 signer,
1578 }
1579 }
1580}
1581
1582fn stream_data_value(spec: &StreamDataSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
1583 let value = Value::Map(base("stream_data", 0, frame_id, sent_at_ms))
1588 .with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
1589 .with_field("seq", Value::Int(spec.seq as i128))
1590 .with_field("encoding", Value::text(spec.encoding.name()))
1591 .with_field("body", spec.body.clone());
1592 with_optional_signer(value, spec.signer)
1593}
1594
1595pub fn stream_data(spec: &StreamDataSpec) -> Value {
1597 stream_data_value(spec, fresh_frame_id(), current_millis())
1598}
1599
1600fn with_optional_signer(value: Value, signer: Option<[u8; 32]>) -> Value {
1604 match signer {
1605 Some(pub_key) => value.with_field("signer", Value::Bytes(pub_key.to_vec())),
1606 None => value,
1607 }
1608}
1609
1610#[derive(Debug, Clone)]
1614pub struct StreamEndSpec {
1615 pub stream_id: [u8; 16],
1616 pub role: StreamRole,
1617 pub signer: Option<[u8; 32]>,
1618}
1619
1620impl StreamEndSpec {
1621 pub fn new(stream_id: [u8; 16], role: StreamRole, signer: Option<[u8; 32]>) -> Self {
1622 Self {
1623 stream_id,
1624 role,
1625 signer,
1626 }
1627 }
1628}
1629
1630fn stream_end_value(spec: &StreamEndSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
1631 let value = Value::Map(base("stream_end", 0, frame_id, sent_at_ms))
1632 .with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
1633 .with_field("role", Value::text(spec.role.name()));
1634 with_optional_signer(value, spec.signer)
1635}
1636
1637pub fn stream_end(spec: &StreamEndSpec) -> Value {
1639 stream_end_value(spec, fresh_frame_id(), current_millis())
1640}
1641
1642#[derive(Debug, Clone)]
1651pub struct StreamErrorSpec {
1652 pub stream_id: [u8; 16],
1653 pub code: String,
1654 pub message: String,
1655 pub signer: Option<[u8; 32]>,
1656}
1657
1658impl StreamErrorSpec {
1659 pub fn new(
1660 stream_id: [u8; 16],
1661 code: impl Into<String>,
1662 message: impl Into<String>,
1663 signer: Option<[u8; 32]>,
1664 ) -> Self {
1665 Self {
1666 stream_id,
1667 code: code.into(),
1668 message: message.into(),
1669 signer,
1670 }
1671 }
1672}
1673
1674fn stream_error_value(spec: &StreamErrorSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
1675 let value = Value::Map(base("stream_error", 0, frame_id, sent_at_ms))
1676 .with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
1677 .with_field("code", Value::Bytes(spec.code.as_bytes().to_vec()))
1678 .with_field("message", Value::Bytes(spec.message.as_bytes().to_vec()));
1679 with_optional_signer(value, spec.signer)
1680}
1681
1682pub fn stream_error(spec: &StreamErrorSpec) -> Value {
1684 stream_error_value(spec, fresh_frame_id(), current_millis())
1685}
1686
1687#[derive(Debug, Clone)]
1691pub struct StreamReplySpec {
1692 pub stream_id: [u8; 16],
1693 pub payload: Value,
1694 pub responded_by: [u8; 32],
1695}
1696
1697impl StreamReplySpec {
1698 pub fn new(stream_id: [u8; 16], payload: Value, responded_by: [u8; 32]) -> Self {
1699 Self {
1700 stream_id,
1701 payload,
1702 responded_by,
1703 }
1704 }
1705}
1706
1707fn stream_reply_value(spec: &StreamReplySpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
1708 Value::Map(base("stream_reply", 0, frame_id, sent_at_ms))
1709 .with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
1710 .with_field("payload", spec.payload.clone())
1711 .with_field("responded_by", Value::Bytes(spec.responded_by.to_vec()))
1712}
1713
1714pub fn stream_reply(spec: &StreamReplySpec) -> Value {
1716 stream_reply_value(spec, fresh_frame_id(), current_millis())
1717}
1718
1719pub fn frame_stream_id(frame: &Value) -> Option<[u8; 16]> {
1724 match frame.get("stream_id") {
1725 Some(Value::Bytes(b)) => b.as_slice().try_into().ok(),
1726 _ => None,
1727 }
1728}
1729
1730#[derive(Debug, Clone)]
1733pub enum StreamEvent {
1734 Data {
1735 stream_id: [u8; 16],
1736 seq: u64,
1737 encoding: StreamEncoding,
1738 body: Value,
1739 },
1740 End {
1741 stream_id: [u8; 16],
1742 role: StreamRole,
1743 },
1744 Error {
1745 stream_id: [u8; 16],
1746 code: String,
1747 message: String,
1748 },
1749 Reply {
1750 stream_id: [u8; 16],
1751 payload: Value,
1752 responded_by: [u8; 32],
1753 },
1754}
1755
1756#[derive(Debug, PartialEq, Eq)]
1757pub enum ParseStreamEventError {
1758 NotAStreamFrame,
1759 MissingField(&'static str),
1760 WrongFieldType(&'static str),
1761}
1762
1763impl std::fmt::Display for ParseStreamEventError {
1764 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1765 match self {
1766 ParseStreamEventError::NotAStreamFrame => write!(
1767 f,
1768 "frame_type is none of stream_data/stream_end/stream_error/stream_reply"
1769 ),
1770 ParseStreamEventError::MissingField(name) => {
1771 write!(f, "missing required field {name:?}")
1772 }
1773 ParseStreamEventError::WrongFieldType(name) => {
1774 write!(f, "field {name:?} has the wrong type")
1775 }
1776 }
1777 }
1778}
1779
1780impl std::error::Error for ParseStreamEventError {}
1781
1782pub fn parse_stream_event(frame: &Value) -> Result<StreamEvent, ParseStreamEventError> {
1785 let stream_id = match frame.get("stream_id") {
1786 Some(Value::Bytes(b)) => b
1787 .as_slice()
1788 .try_into()
1789 .map_err(|_| ParseStreamEventError::WrongFieldType("stream_id"))?,
1790 Some(_) => return Err(ParseStreamEventError::WrongFieldType("stream_id")),
1791 None => return Err(ParseStreamEventError::MissingField("stream_id")),
1792 };
1793 match frame.get("frame_type") {
1794 Some(Value::Text(t)) if t == "stream_data" => {
1795 let seq = match frame.get("seq") {
1796 Some(Value::Int(n)) if *n >= 0 => *n as u64,
1797 Some(_) => return Err(ParseStreamEventError::WrongFieldType("seq")),
1798 None => return Err(ParseStreamEventError::MissingField("seq")),
1799 };
1800 let encoding = match frame.get("encoding") {
1801 Some(Value::Text(t)) => StreamEncoding::from_name(t)
1802 .ok_or(ParseStreamEventError::WrongFieldType("encoding"))?,
1803 Some(_) => return Err(ParseStreamEventError::WrongFieldType("encoding")),
1804 None => return Err(ParseStreamEventError::MissingField("encoding")),
1805 };
1806 let body = frame
1807 .get("body")
1808 .cloned()
1809 .ok_or(ParseStreamEventError::MissingField("body"))?;
1810 Ok(StreamEvent::Data {
1811 stream_id,
1812 seq,
1813 encoding,
1814 body,
1815 })
1816 }
1817 Some(Value::Text(t)) if t == "stream_end" => {
1818 let role = match frame.get("role") {
1819 Some(Value::Text(t)) => {
1820 StreamRole::from_name(t).ok_or(ParseStreamEventError::WrongFieldType("role"))?
1821 }
1822 Some(_) => return Err(ParseStreamEventError::WrongFieldType("role")),
1823 None => return Err(ParseStreamEventError::MissingField("role")),
1824 };
1825 Ok(StreamEvent::End { stream_id, role })
1826 }
1827 Some(Value::Text(t)) if t == "stream_error" => {
1828 let code = match frame.get("code") {
1829 Some(Value::Bytes(b)) => String::from_utf8(b.clone())
1830 .map_err(|_| ParseStreamEventError::WrongFieldType("code"))?,
1831 Some(_) => return Err(ParseStreamEventError::WrongFieldType("code")),
1832 None => return Err(ParseStreamEventError::MissingField("code")),
1833 };
1834 let message = match frame.get("message") {
1835 Some(Value::Bytes(b)) => String::from_utf8(b.clone())
1836 .map_err(|_| ParseStreamEventError::WrongFieldType("message"))?,
1837 Some(_) => return Err(ParseStreamEventError::WrongFieldType("message")),
1838 None => return Err(ParseStreamEventError::MissingField("message")),
1839 };
1840 Ok(StreamEvent::Error {
1841 stream_id,
1842 code,
1843 message,
1844 })
1845 }
1846 Some(Value::Text(t)) if t == "stream_reply" => {
1847 let payload = frame
1848 .get("payload")
1849 .cloned()
1850 .ok_or(ParseStreamEventError::MissingField("payload"))?;
1851 let responded_by = match frame.get("responded_by") {
1852 Some(Value::Bytes(b)) => b
1853 .as_slice()
1854 .try_into()
1855 .map_err(|_| ParseStreamEventError::WrongFieldType("responded_by"))?,
1856 Some(_) => return Err(ParseStreamEventError::WrongFieldType("responded_by")),
1857 None => return Err(ParseStreamEventError::MissingField("responded_by")),
1858 };
1859 Ok(StreamEvent::Reply {
1860 stream_id,
1861 payload,
1862 responded_by,
1863 })
1864 }
1865 Some(_) | None => Err(ParseStreamEventError::NotAStreamFrame),
1866 }
1867}
1868
1869#[cfg(test)]
1870mod tests {
1871 use super::*;
1872
1873 fn hex_bytes(s: &str) -> Vec<u8> {
1874 ::hex::decode(s).expect("valid hex fixture")
1875 }
1876
1877 fn fixed_array(hex_str: &str) -> [u8; 32] {
1878 hex_bytes(hex_str).try_into().expect("32-byte fixture")
1879 }
1880
1881 const VECTOR_PUB: &str = "B966A9812649C3D5542FF54954FE090C43FDA6574FE48A0DD326626CFAD29A83";
1884 const VECTOR_PRIV: &str = "457F45FF5A09E172ED15CB20D6CB26B51AD15ED7308C12D478E8631F9CA03D4F";
1885 const VECTOR_PUZZLE_EVIDENCE: &str =
1886 "09D48C91CB46513ED2580BDCEA87C40DA508D4E50EC3DF2F701AFC55D1C5C0B2";
1887 const VECTOR_FRAME_ID: &str = "0192E8B0F1A47000A1B2C3D4E5F60718";
1888 const VECTOR_SENT_AT_MS: u64 = 1_700_000_000_000;
1889 const VECTOR_SIGNATURE: &str = "CF6959A61A2F4D2046F0124C1DD56A6541265F36A24CB18CA8C45C95031854D6AECE5FB93E2AE7BA6C444A09C7C5DED195B6EB0D1CC8E487CCF6E4F0D903B409";
1890 const VECTOR_ENCODED_LEN: usize = 375;
1891
1892 #[test]
1900 fn connect_frame_matches_the_reference_byte_for_byte() {
1901 let pub_bytes = fixed_array(VECTOR_PUB);
1902 let identity = KeyPair::from_seed_bytes(fixed_array(VECTOR_PRIV));
1903 let puzzle_evidence = fixed_array(VECTOR_PUZZLE_EVIDENCE);
1904 let frame_id: [u8; 16] = hex_bytes(VECTOR_FRAME_ID).try_into().expect("16 bytes");
1905
1906 let spec = ConnectSpec::new(pub_bytes, puzzle_evidence);
1907 let unsigned = connect_value(&spec, frame_id, VECTOR_SENT_AT_MS);
1908 let signed = sign(unsigned, &identity);
1909
1910 let sig_field = match signed.get("signature") {
1911 Some(Value::Bytes(b)) => b.clone(),
1912 other => panic!("expected a signature field, got {other:?}"),
1913 };
1914 assert_eq!(
1915 hex::encode_upper(&sig_field),
1916 VECTOR_SIGNATURE,
1917 "signature diverged from the reference — canonical CBOR encoding \
1918 or the signing domain/bytes must differ somewhere"
1919 );
1920
1921 let encoded = encode(&signed).expect("encodable frame");
1922 assert_eq!(encoded.len(), VECTOR_ENCODED_LEN);
1923
1924 let decoded = match decode(&encoded).expect("valid frame") {
1927 Decoded::Frame(value, consumed) => {
1928 assert_eq!(consumed, encoded.len());
1929 value
1930 }
1931 Decoded::More(n) => panic!("unexpectedly needed {n} more bytes"),
1932 };
1933 verify(&decoded, &pub_bytes).expect("our own signature must verify");
1934 }
1935
1936 #[test]
1937 fn verify_rejects_a_tampered_field() {
1938 let identity = KeyPair::from_seed_bytes(fixed_array(VECTOR_PRIV));
1939 let pub_bytes = identity.public_bytes();
1940 let spec = ConnectSpec::new(pub_bytes, fixed_array(VECTOR_PUZZLE_EVIDENCE));
1941 let signed = sign(connect(&spec), &identity);
1942
1943 let tampered = signed.with_field("capabilities", Value::Int(999));
1945 assert_eq!(
1946 verify(&tampered, &pub_bytes),
1947 Err(VerifyError::SignatureInvalid)
1948 );
1949 }
1950
1951 #[test]
1952 fn verify_rejects_a_missing_signature() {
1953 let frame = Value::Map(vec![(Value::text("frame_type"), Value::text("connect"))]);
1954 let pubkey = [0u8; 32];
1955 assert_eq!(verify(&frame, &pubkey), Err(VerifyError::MissingSignature));
1956 }
1957
1958 #[test]
1959 fn decode_reports_more_for_a_short_buffer() {
1960 assert!(matches!(decode(&[0, 0]), Ok(Decoded::More(2))));
1961 let mut buf = 10u32.to_be_bytes().to_vec();
1964 buf.extend_from_slice(&[0, 0]);
1965 assert!(matches!(decode(&buf), Ok(Decoded::More(8))));
1966 }
1967
1968 #[test]
1969 fn decode_rejects_a_length_over_the_cap() {
1970 let buf = ((MAX_FRAME_BYTES as u32) + 1).to_be_bytes();
1971 assert!(matches!(
1972 decode(&buf),
1973 Err(DecodeFrameError::TooLarge(n)) if n == MAX_FRAME_BYTES + 1
1974 ));
1975 }
1976
1977 #[test]
1978 fn goodbye_frame_round_trips() {
1979 let frame = goodbye("normal", Some("bye"));
1980 assert_eq!(frame.get("frame_type"), Some(&Value::text("goodbye")));
1981 assert_eq!(frame.get("reason"), Some(&Value::text("normal")));
1982 assert_eq!(frame.get("detail"), Some(&Value::Bytes(b"bye".to_vec())));
1983 }
1984
1985 #[test]
1986 fn goodbye_without_detail_is_null() {
1987 let frame = goodbye("timeout", None);
1988 assert_eq!(frame.get("detail"), Some(&Value::Null));
1989 }
1990
1991 #[test]
1992 fn parse_hello_reads_a_well_formed_frame() {
1993 let node_id = [7u8; 32];
1994 let station_id = [8u8; 32];
1995 let realm = [9u8; 32];
1996 let hello = Value::Map(vec![
1997 (Value::text("frame_type"), Value::text("hello")),
1998 (Value::text("node_id"), Value::Bytes(node_id.to_vec())),
1999 (Value::text("station_id"), Value::Bytes(station_id.to_vec())),
2000 (
2001 Value::text("realms"),
2002 Value::List(vec![Value::Bytes(realm.to_vec())]),
2003 ),
2004 (Value::text("capabilities"), Value::Int(0)),
2005 (Value::text("accepted"), Value::text("true")),
2006 (Value::text("negotiated_capabilities"), Value::Int(3)),
2007 ]);
2008 let info = parse_hello(&hello).expect("well-formed hello");
2009 assert_eq!(info.node_id, node_id);
2010 assert_eq!(info.station_id, station_id);
2011 assert_eq!(info.realms, vec![realm]);
2012 assert!(info.accepted);
2013 assert_eq!(info.negotiated_capabilities, 3);
2014 assert_eq!(info.refusal_code, None);
2015 }
2016
2017 #[test]
2018 fn parse_hello_rejects_the_wrong_frame_type() {
2019 let frame = Value::Map(vec![(Value::text("frame_type"), Value::text("connect"))]);
2020 assert_eq!(parse_hello(&frame), Err(ParseHelloError::NotAHelloFrame));
2021 }
2022
2023 #[test]
2024 fn parse_hello_reports_a_missing_field() {
2025 let frame = Value::Map(vec![(Value::text("frame_type"), Value::text("hello"))]);
2026 assert_eq!(
2027 parse_hello(&frame),
2028 Err(ParseHelloError::MissingField("node_id"))
2029 );
2030 }
2031
2032 const VECTOR_CALL_ID: &str = "AABBCCDDEEFF00112233445566778899";
2045 const VECTOR_ZERO_REALM: [u8; 32] = [0u8; 32];
2046
2047 fn vector_identity() -> KeyPair {
2048 KeyPair::from_seed_bytes(fixed_array(VECTOR_PRIV))
2049 }
2050
2051 fn vector_call_id() -> [u8; 16] {
2052 hex_bytes(VECTOR_CALL_ID).try_into().expect("16 bytes")
2053 }
2054
2055 fn vector_frame_id() -> [u8; 16] {
2056 hex_bytes(VECTOR_FRAME_ID).try_into().expect("16 bytes")
2057 }
2058
2059 #[test]
2060 fn call_frame_matches_the_reference_byte_for_byte() {
2061 let pub_bytes = fixed_array(VECTOR_PUB);
2062 let identity = vector_identity();
2063 let spec = CallSpec::new(
2064 vector_call_id(),
2065 "_content.get_manifest",
2066 VECTOR_ZERO_REALM,
2067 Value::Map(vec![(Value::text("hello"), Value::text("world"))]),
2068 1_700_000_030_000,
2069 pub_bytes,
2070 );
2071 let signed = sign(
2072 call_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2073 &identity,
2074 );
2075 let sig = match signed.get("signature") {
2076 Some(Value::Bytes(b)) => b.clone(),
2077 other => panic!("expected a signature field, got {other:?}"),
2078 };
2079 assert_eq!(
2080 hex::encode_upper(&sig),
2081 "A6BC174F0241E644F634702C08781C8FC8BD3CDE3CA9650DE8A731A01203D9B9403A2CAD75800F7B8C9AAE16FA146B1195FF03F0E6DC4595A652D7F29BFE350A"
2082 );
2083 let encoded = encode(&signed).expect("encodable frame");
2084 assert_eq!(encoded.len(), 386);
2085 }
2086
2087 #[test]
2088 fn result_frame_matches_the_reference_byte_for_byte() {
2089 let pub_bytes = fixed_array(VECTOR_PUB);
2090 let identity = vector_identity();
2091 let spec = ResultSpec::new(vector_call_id(), Value::text("ok-result"), pub_bytes);
2092 let signed = sign(
2093 result_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2094 &identity,
2095 );
2096 let sig = match signed.get("signature") {
2097 Some(Value::Bytes(b)) => b.clone(),
2098 other => panic!("expected a signature field, got {other:?}"),
2099 };
2100 assert_eq!(
2101 hex::encode_upper(&sig),
2102 "03E8F72D51D958C318B7F1C25D78408408317DEAB23434D6EA32F211CADEA1C62900DA15AFF603E795B19A388D382BDB10E65AEFC6F0CE551270AB172A88E50B"
2103 );
2104 assert_eq!(encode(&signed).expect("encodable").len(), 301);
2105 }
2106
2107 #[test]
2108 fn error_frame_matches_the_reference_byte_for_byte() {
2109 let pub_bytes = fixed_array(VECTOR_PUB);
2110 let identity = vector_identity();
2111 let spec = CallErrorSpec::new(
2112 vector_call_id(),
2113 crate::bolt4::Code::UnknownNextPeer,
2114 pub_bytes,
2115 );
2116 let signed = sign(
2117 call_error_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2118 &identity,
2119 );
2120 let sig = match signed.get("signature") {
2121 Some(Value::Bytes(b)) => b.clone(),
2122 other => panic!("expected a signature field, got {other:?}"),
2123 };
2124 assert_eq!(
2125 hex::encode_upper(&sig),
2126 "182ECD5217CE378F576635B23CC8C9F265555142845D6CBA033A282BAED97966C23FBE91D08507FB8E840375AA17665763804F40F89102F8D3EDAD4DA98FC20D"
2127 );
2128 assert_eq!(encode(&signed).expect("encodable").len(), 333);
2129 }
2130
2131 #[test]
2132 fn publish_frame_matches_the_reference_byte_for_byte() {
2133 let pub_bytes = fixed_array(VECTOR_PUB);
2134 let identity = vector_identity();
2135 let spec = PublishSpec::new(
2136 "test.topic",
2137 VECTOR_ZERO_REALM,
2138 pub_bytes,
2139 42,
2140 Value::text("published-data"),
2141 VECTOR_SENT_AT_MS,
2142 );
2143 let signed = sign(
2144 publish_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2145 &identity,
2146 );
2147 let sig = match signed.get("signature") {
2148 Some(Value::Bytes(b)) => b.clone(),
2149 other => panic!("expected a signature field, got {other:?}"),
2150 };
2151 assert_eq!(
2152 hex::encode_upper(&sig),
2153 "DD49D10EFA9F2EED0A393DC02DC5BBAC25D6731562EA39F5AB2E5337824527AFFBC7D917AF4DE5EFDBE5BC41E58659E05EC6FDE4E91FB1A32CC9C211456DF10C"
2154 );
2155 assert_eq!(encode(&signed).expect("encodable").len(), 355);
2156 }
2157
2158 #[test]
2167 fn publisher_sig_matches_the_erlang_reference() {
2168 let pub_bytes = fixed_array(VECTOR_PUB);
2169 let identity = vector_identity();
2170 let spec = PublishSpec::new(
2171 "acme/svc.do",
2172 VECTOR_ZERO_REALM,
2173 pub_bytes,
2174 42,
2175 Value::Bytes(b"hello".to_vec()),
2176 VECTOR_SENT_AT_MS,
2177 );
2178 let unsigned = publish_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS);
2179 let with_pub_sig = sign_publisher(unsigned, &identity);
2180
2181 let sig = match with_pub_sig.get("publisher_sig") {
2182 Some(Value::Bytes(b)) => b.clone(),
2183 other => panic!("expected a publisher_sig field, got {other:?}"),
2184 };
2185 assert_eq!(
2186 hex::encode_upper(&sig),
2187 "C11BEB676A590FD1BA86F0B77E377B4582AA461DB1283F64E57224E920A7BD0A2C7D36271B795FFC3CB4F2C7BB8925B034431AA6425E25B2AEEFAC026883BB0C"
2188 );
2189
2190 verify_publisher(&with_pub_sig).expect("our own freshly-signed frame must verify");
2191
2192 let tampered = with_pub_sig
2194 .clone()
2195 .with_field("payload", Value::Bytes(b"world".to_vec()));
2196 assert!(
2197 verify_publisher(&tampered).is_err(),
2198 "verify_publisher accepted a frame with a tampered payload"
2199 );
2200
2201 assert_eq!(
2203 verify_publisher(&unsigned_publish_for_tamper_check(&spec)),
2204 Err(VerifyPublisherError::MissingPublisherSig)
2205 );
2206 }
2207
2208 fn unsigned_publish_for_tamper_check(spec: &PublishSpec) -> Value {
2209 publish_value(spec, vector_frame_id(), VECTOR_SENT_AT_MS)
2210 }
2211
2212 #[test]
2217 fn publish_frame_with_both_signatures_round_trips() {
2218 let identity = KeyPair::generate();
2219 let pub_bytes = identity.node_id();
2220 let spec = PublishSpec::new(
2221 "acme/svc.do",
2222 VECTOR_ZERO_REALM,
2223 pub_bytes,
2224 1,
2225 Value::Bytes(b"hello".to_vec()),
2226 VECTOR_SENT_AT_MS,
2227 );
2228 let unsigned = publish(&spec);
2229 let with_pub_sig = sign_publisher(unsigned, &identity);
2230 let fully_signed = sign(with_pub_sig, &identity);
2231
2232 let encoded = encode(&fully_signed).expect("encodable");
2233 let decoded = match decode(&encoded).expect("decodable") {
2234 Decoded::Frame(value, consumed) => {
2235 assert_eq!(consumed, encoded.len());
2236 value
2237 }
2238 Decoded::More(n) => panic!("unexpectedly needed {n} more bytes"),
2239 };
2240
2241 verify(&decoded, &pub_bytes).expect("per-hop verify on decoded frame");
2242 verify_publisher(&decoded).expect("verify_publisher on decoded frame");
2243
2244 assert!(
2245 decoded.get("publisher_sig").is_some(),
2246 "decoded frame lost publisher_sig"
2247 );
2248 assert!(
2249 decoded.get("signature").is_some(),
2250 "decoded frame lost signature"
2251 );
2252 }
2253
2254 #[test]
2255 fn subscribe_frame_matches_the_reference_byte_for_byte() {
2256 let pub_bytes = fixed_array(VECTOR_PUB);
2257 let identity = vector_identity();
2258 let spec = SubscribeSpec::new("test.topic", VECTOR_ZERO_REALM, pub_bytes);
2259 let signed = sign(
2260 subscribe_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2261 &identity,
2262 );
2263 let sig = match signed.get("signature") {
2264 Some(Value::Bytes(b)) => b.clone(),
2265 other => panic!("expected a signature field, got {other:?}"),
2266 };
2267 assert_eq!(
2268 hex::encode_upper(&sig),
2269 "ABDD7304B887A53B149CE4D4C62F1AFD20AE07D8612B76F22006FA6676B8DDB37C1D5106358D32080246BA4355A9E04BF49F73600E752F5F9037D7A93A47020A"
2270 );
2271 assert_eq!(encode(&signed).expect("encodable").len(), 313);
2272 }
2273
2274 #[test]
2275 fn unsubscribe_frame_matches_the_reference_byte_for_byte() {
2276 let pub_bytes = fixed_array(VECTOR_PUB);
2277 let identity = vector_identity();
2278 let spec = UnsubscribeSpec::new("test.topic", VECTOR_ZERO_REALM, pub_bytes);
2279 let signed = sign(
2280 unsubscribe_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2281 &identity,
2282 );
2283 let sig = match signed.get("signature") {
2284 Some(Value::Bytes(b)) => b.clone(),
2285 other => panic!("expected a signature field, got {other:?}"),
2286 };
2287 assert_eq!(
2288 hex::encode_upper(&sig),
2289 "C917068BE4E1C5A3C753F249037DD8F44293D888BB252BF1E828671969547969982160C91A0E3CA1C31DE29ED39E3677E7F20F4BDE61539D4618B3703018E403"
2290 );
2291 assert_eq!(encode(&signed).expect("encodable").len(), 298);
2292 }
2293
2294 #[test]
2295 fn event_frame_matches_the_reference_byte_for_byte() {
2296 let pub_bytes = fixed_array(VECTOR_PUB);
2297 let identity = vector_identity();
2298 let fields = base("event", 0, vector_frame_id(), VECTOR_SENT_AT_MS);
2299 let unsigned = Value::Map(fields)
2300 .with_field("realm", Value::Bytes(VECTOR_ZERO_REALM.to_vec()))
2301 .with_field("topic", Value::Bytes(b"test.topic".to_vec()))
2302 .with_field("publisher", Value::Bytes(pub_bytes.to_vec()))
2303 .with_field("seq", Value::Int(42))
2304 .with_field("payload", Value::text("published-data"))
2305 .with_field("delivered_via", Value::text("direct"));
2306 let signed = sign(unsigned, &identity);
2307 let sig = match signed.get("signature") {
2308 Some(Value::Bytes(b)) => b.clone(),
2309 other => panic!("expected a signature field, got {other:?}"),
2310 };
2311 assert_eq!(
2312 hex::encode_upper(&sig),
2313 "9B9EE4EAC375FBD0C9B5A5BC6D82E35739F8ECBF594979891BF35E5BDB53A148B3936AF99217C3D8C12E2EEA0686F68D5FE63284BE6B142F87BFF319DDDB780F"
2314 );
2315 assert_eq!(encode(&signed).expect("encodable").len(), 341);
2316
2317 let decoded = decode(&encode(&signed).unwrap()).unwrap();
2320 let Decoded::Frame(value, _) = decoded else {
2321 panic!("expected a complete frame")
2322 };
2323 let info = parse_event(&value).expect("well-formed event");
2324 assert_eq!(info.topic, "test.topic");
2325 assert_eq!(info.seq, 42);
2326 assert_eq!(info.delivered_via, "direct");
2327 }
2328
2329 #[test]
2330 fn parse_call_response_reads_a_result() {
2331 let frame = Value::Map(vec![
2332 (Value::text("frame_type"), Value::text("result")),
2333 (Value::text("call_id"), Value::Bytes(vec![1; 16])),
2334 (Value::text("payload"), Value::text("ok")),
2335 (Value::text("responded_by"), Value::Bytes(vec![2; 32])),
2336 ]);
2337 match parse_call_response(&frame).expect("well-formed result") {
2338 CallResponse::Result {
2339 payload,
2340 responded_by,
2341 } => {
2342 assert_eq!(payload, Value::text("ok"));
2343 assert_eq!(responded_by, [2u8; 32]);
2344 }
2345 other => panic!("expected Result, got {other:?}"),
2346 }
2347 }
2348
2349 #[test]
2350 fn parse_call_response_reads_an_error() {
2351 let frame = Value::Map(vec![
2352 (Value::text("frame_type"), Value::text("error")),
2353 (Value::text("call_id"), Value::Bytes(vec![1; 16])),
2354 (Value::text("code"), Value::Int(1)),
2355 (Value::text("name"), Value::text("unknown_next_peer")),
2356 (Value::text("reported_by"), Value::Bytes(vec![2; 32])),
2357 (Value::text("detail"), Value::Null),
2358 ]);
2359 match parse_call_response(&frame).expect("well-formed error") {
2360 CallResponse::Error {
2361 code,
2362 name,
2363 reported_by,
2364 detail,
2365 } => {
2366 assert_eq!(code, 1);
2367 assert_eq!(name, "unknown_next_peer");
2368 assert_eq!(reported_by, [2u8; 32]);
2369 assert_eq!(detail, None);
2370 }
2371 other => panic!("expected Error, got {other:?}"),
2372 }
2373 }
2374
2375 #[test]
2376 fn frame_call_id_reads_from_any_frame_type() {
2377 let frame = Value::Map(vec![(Value::text("call_id"), Value::Bytes(vec![9; 16]))]);
2378 assert_eq!(frame_call_id(&frame), Some([9u8; 16]));
2379 let wrong_size = Value::Map(vec![(Value::text("call_id"), Value::Bytes(vec![9; 32]))]);
2382 assert_eq!(frame_call_id(&wrong_size), None);
2383 }
2384
2385 const VECTOR_STREAM_ID: &str = "0102030405060708090A0B0C0D0E0F10";
2395
2396 fn vector_stream_id() -> [u8; 16] {
2397 hex_bytes(VECTOR_STREAM_ID).try_into().expect("16 bytes")
2398 }
2399
2400 #[test]
2401 fn stream_open_frame_matches_the_reference_byte_for_byte() {
2402 let pub_bytes = fixed_array(VECTOR_PUB);
2403 let identity = vector_identity();
2404 let spec = StreamOpenSpec::new(
2405 vector_stream_id(),
2406 "macula_rust_sdk.test_stream",
2407 VECTOR_ZERO_REALM,
2408 StreamMode::ClientStream,
2409 Value::Map(vec![(Value::text("hello"), Value::text("world"))]),
2410 1_700_000_030_000,
2411 pub_bytes,
2412 );
2413 let signed = sign(
2414 stream_open_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2415 &identity,
2416 );
2417 let sig = match signed.get("signature") {
2418 Some(Value::Bytes(b)) => b.clone(),
2419 other => panic!("expected a signature field, got {other:?}"),
2420 };
2421 assert_eq!(hex::encode_upper(&sig), "6070D8AB71F837591AC2C803C04F9E1D3FA01C9310D33C96A90434820C5E50550F9DEA8A764247EB49AF63447C037E192B7892A365C1A4ACB9BC46B98AA5670F");
2422 let encoded = encode(&signed).expect("encodable frame");
2423 assert_eq!(encoded.len(), 415);
2424 }
2425
2426 #[test]
2433 fn parse_stream_open_round_trips_a_well_formed_frame() {
2434 let pub_bytes = fixed_array(VECTOR_PUB);
2435 let spec = StreamOpenSpec::new(
2436 vector_stream_id(),
2437 "macula_rust_sdk.test_stream",
2438 VECTOR_ZERO_REALM,
2439 StreamMode::ClientStream,
2440 Value::Map(vec![(Value::text("hello"), Value::text("world"))]),
2441 1_700_000_030_000,
2442 pub_bytes,
2443 );
2444 let frame = stream_open_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS);
2445 let info = parse_stream_open(&frame).expect("well-formed stream_open");
2446 assert_eq!(info.stream_id, vector_stream_id());
2447 assert_eq!(info.procedure, "macula_rust_sdk.test_stream");
2448 assert_eq!(info.realm, VECTOR_ZERO_REALM);
2449 assert_eq!(info.mode, StreamMode::ClientStream);
2450 assert_eq!(
2451 info.args,
2452 Value::Map(vec![(Value::text("hello"), Value::text("world"))])
2453 );
2454 assert_eq!(info.deadline_ms, 1_700_000_030_000);
2455 assert_eq!(info.caller, pub_bytes);
2456 }
2457
2458 #[test]
2459 fn parse_stream_open_rejects_the_wrong_frame_type() {
2460 let frame = Value::Map(vec![(
2461 Value::text("frame_type"),
2462 Value::text("stream_data"),
2463 )]);
2464 assert_eq!(
2465 parse_stream_open(&frame).unwrap_err(),
2466 ParseStreamOpenError::NotAStreamOpenFrame
2467 );
2468 }
2469
2470 #[test]
2471 fn stream_data_raw_frame_matches_the_reference_byte_for_byte() {
2472 let identity = vector_identity();
2473 let spec = StreamDataSpec::new(
2474 vector_stream_id(),
2475 0,
2476 StreamEncoding::Raw,
2477 Value::Bytes(b"raw chunk bytes".to_vec()),
2478 None,
2479 );
2480 let signed = sign(
2481 stream_data_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2482 &identity,
2483 );
2484 let sig = match signed.get("signature") {
2485 Some(Value::Bytes(b)) => b.clone(),
2486 other => panic!("expected a signature field, got {other:?}"),
2487 };
2488 assert_eq!(hex::encode_upper(&sig), "35770744FE5BD01B86DDA01AB4EF855E4E4FE0EDFEDC89FF690728C585C60A5CB035717E3EA9133C4AD833E226F4DB95E9A5AF9AC59E7BACBB8BDF72611F8003");
2489 let encoded = encode(&signed).expect("encodable frame");
2490 assert_eq!(encoded.len(), 269);
2491 }
2492
2493 #[test]
2500 fn stream_data_with_signer_matches_the_reference_byte_for_byte() {
2501 let identity = vector_identity();
2502 let spec = StreamDataSpec::new(
2503 vector_stream_id(),
2504 0,
2505 StreamEncoding::Raw,
2506 Value::Bytes(b"raw chunk bytes".to_vec()),
2507 Some(fixed_array(VECTOR_PUB)),
2508 );
2509 let signed = sign(
2510 stream_data_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2511 &identity,
2512 );
2513 let sig = match signed.get("signature") {
2514 Some(Value::Bytes(b)) => b.clone(),
2515 other => panic!("expected a signature field, got {other:?}"),
2516 };
2517 assert_eq!(hex::encode_upper(&sig), "3EA0B6B6DB1549D2EA42AF015A477FCD6D00B11F48F9CC07AF0914CAC18F22B5C12E5EE446811388F207D688960B67D9BEE7B4D998BE02F2B1426B6C4A06D307");
2518 let encoded = encode(&signed).expect("encodable frame");
2519 assert_eq!(encoded.len(), 310);
2520 }
2521
2522 #[test]
2530 fn stream_data_msgpack_frame_matches_the_reference_byte_for_byte() {
2531 let identity = vector_identity();
2532 let spec = StreamDataSpec::new(
2533 vector_stream_id(),
2534 1,
2535 StreamEncoding::Msgpack,
2536 Value::Map(vec![
2537 (Value::text("a"), Value::Int(1)),
2538 (Value::text("greeting"), Value::Bytes(b"hi".to_vec())),
2542 ]),
2543 None,
2544 );
2545 let signed = sign(
2546 stream_data_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2547 &identity,
2548 );
2549 let sig = match signed.get("signature") {
2550 Some(Value::Bytes(b)) => b.clone(),
2551 other => panic!("expected a signature field, got {other:?}"),
2552 };
2553 assert_eq!(hex::encode_upper(&sig), "99CA90B0C01FD349DBAF317D03872E5F460426789874D79B6FBE37F4AC92C2AD690A00CDB3734F262D5C58C8F3BFD06F8AE892A8B5655274718A283ABA1D4D08");
2554 let encoded = encode(&signed).expect("encodable frame");
2555 assert_eq!(encoded.len(), 273);
2556 }
2557
2558 #[test]
2559 fn stream_end_frame_matches_the_reference_byte_for_byte() {
2560 let identity = vector_identity();
2561 let spec = StreamEndSpec::new(vector_stream_id(), StreamRole::Send, None);
2562 let signed = sign(
2563 stream_end_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2564 &identity,
2565 );
2566 let sig = match signed.get("signature") {
2567 Some(Value::Bytes(b)) => b.clone(),
2568 other => panic!("expected a signature field, got {other:?}"),
2569 };
2570 assert_eq!(hex::encode_upper(&sig), "78F2B94BD5AC70901EABB31D8B17C89B58A88942300C6232545899AFB933B2C4B7399BB183A5660671981B6346DA27033C8F93A99E7EBA96F0F689B03D4F940A");
2571 let encoded = encode(&signed).expect("encodable frame");
2572 assert_eq!(encoded.len(), 239);
2573 }
2574
2575 #[test]
2578 fn stream_end_with_signer_matches_the_reference_byte_for_byte() {
2579 let identity = vector_identity();
2580 let spec = StreamEndSpec::new(
2581 vector_stream_id(),
2582 StreamRole::Send,
2583 Some(fixed_array(VECTOR_PUB)),
2584 );
2585 let signed = sign(
2586 stream_end_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2587 &identity,
2588 );
2589 let sig = match signed.get("signature") {
2590 Some(Value::Bytes(b)) => b.clone(),
2591 other => panic!("expected a signature field, got {other:?}"),
2592 };
2593 assert_eq!(hex::encode_upper(&sig), "CC316B0A1C1AD4701AD16D8A140ED62D5DEEFD721C1CEB574CC8755C645CA27413EF9C6A6A9C4768564524C412515C14637A9D6BD4CCB8CD1ADD44F2A240C70C");
2594 let encoded = encode(&signed).expect("encodable frame");
2595 assert_eq!(encoded.len(), 280);
2596 }
2597
2598 #[test]
2599 fn stream_error_frame_matches_the_reference_byte_for_byte() {
2600 let identity = vector_identity();
2601 let spec = StreamErrorSpec::new(vector_stream_id(), "cancelled", "boom", None);
2602 let signed = sign(
2603 stream_error_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2604 &identity,
2605 );
2606 let sig = match signed.get("signature") {
2607 Some(Value::Bytes(b)) => b.clone(),
2608 other => panic!("expected a signature field, got {other:?}"),
2609 };
2610 assert_eq!(hex::encode_upper(&sig), "119F379518EC17C603ED5466A57D7AE53198A8AC4D5CA9849934A78994428CB3DAD40BC0EFECE1A0C8EEB0ACC28973C0F7E55DE6444827091814AF0715D9FF0B");
2611 let encoded = encode(&signed).expect("encodable frame");
2612 assert_eq!(encoded.len(), 259);
2613 }
2614
2615 #[test]
2618 fn stream_error_with_signer_matches_the_reference_byte_for_byte() {
2619 let identity = vector_identity();
2620 let spec = StreamErrorSpec::new(
2621 vector_stream_id(),
2622 "cancelled",
2623 "boom",
2624 Some(fixed_array(VECTOR_PUB)),
2625 );
2626 let signed = sign(
2627 stream_error_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2628 &identity,
2629 );
2630 let sig = match signed.get("signature") {
2631 Some(Value::Bytes(b)) => b.clone(),
2632 other => panic!("expected a signature field, got {other:?}"),
2633 };
2634 assert_eq!(hex::encode_upper(&sig), "223062E2816C5E6DABCF08A0A4FD01F477F2D1D933F2F1FDC971CAB570003DDE8192CC2F8811CE4A2D180B6781AFA64EB4057947E25CF121F745A9654DC23D0A");
2635 let encoded = encode(&signed).expect("encodable frame");
2636 assert_eq!(encoded.len(), 300);
2637 }
2638
2639 #[test]
2640 fn stream_reply_frame_matches_the_reference_byte_for_byte() {
2641 let pub_bytes = fixed_array(VECTOR_PUB);
2642 let identity = vector_identity();
2643 let spec = StreamReplySpec::new(
2644 vector_stream_id(),
2645 Value::Map(vec![(Value::text("ok"), Value::text("true"))]),
2646 pub_bytes,
2647 );
2648 let signed = sign(
2649 stream_reply_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2650 &identity,
2651 );
2652 let sig = match signed.get("signature") {
2653 Some(Value::Bytes(b)) => b.clone(),
2654 other => panic!("expected a signature field, got {other:?}"),
2655 };
2656 assert_eq!(hex::encode_upper(&sig), "ADF57AD58B253F175ADF72E4717E078C62F3E22CBDDBF8DDC0DD8A47CAAA061E8A37C73BAAB91E450D1D8472021B6A0161169D77E9D186C436D3E6580D48C703");
2657 let encoded = encode(&signed).expect("encodable frame");
2658 assert_eq!(encoded.len(), 295);
2659 }
2660
2661 #[test]
2662 fn frame_stream_id_reads_from_any_frame_type() {
2663 let frame = Value::Map(vec![(Value::text("stream_id"), Value::Bytes(vec![9; 16]))]);
2664 assert_eq!(frame_stream_id(&frame), Some([9u8; 16]));
2665 let wrong_size = Value::Map(vec![(Value::text("stream_id"), Value::Bytes(vec![9; 32]))]);
2666 assert_eq!(frame_stream_id(&wrong_size), None);
2667 }
2668
2669 #[test]
2670 fn parse_stream_event_reads_data_end_error_and_reply() {
2671 let data = Value::Map(vec![
2672 (Value::text("frame_type"), Value::text("stream_data")),
2673 (Value::text("stream_id"), Value::Bytes(vec![1; 16])),
2674 (Value::text("seq"), Value::Int(3)),
2675 (Value::text("encoding"), Value::text("raw")),
2676 (Value::text("body"), Value::Bytes(b"hi".to_vec())),
2677 ]);
2678 match parse_stream_event(&data).expect("well-formed stream_data") {
2679 StreamEvent::Data {
2680 stream_id,
2681 seq,
2682 encoding,
2683 body,
2684 } => {
2685 assert_eq!(stream_id, [1u8; 16]);
2686 assert_eq!(seq, 3);
2687 assert_eq!(encoding, StreamEncoding::Raw);
2688 assert_eq!(body, Value::Bytes(b"hi".to_vec()));
2689 }
2690 other => panic!("expected Data, got {other:?}"),
2691 }
2692
2693 let end = Value::Map(vec![
2694 (Value::text("frame_type"), Value::text("stream_end")),
2695 (Value::text("stream_id"), Value::Bytes(vec![1; 16])),
2696 (Value::text("role"), Value::text("both")),
2697 ]);
2698 match parse_stream_event(&end).expect("well-formed stream_end") {
2699 StreamEvent::End { stream_id, role } => {
2700 assert_eq!(stream_id, [1u8; 16]);
2701 assert_eq!(role, StreamRole::Both);
2702 }
2703 other => panic!("expected End, got {other:?}"),
2704 }
2705
2706 let error = Value::Map(vec![
2707 (Value::text("frame_type"), Value::text("stream_error")),
2708 (Value::text("stream_id"), Value::Bytes(vec![1; 16])),
2709 (Value::text("code"), Value::Bytes(b"cancelled".to_vec())),
2710 (Value::text("message"), Value::Bytes(b"boom".to_vec())),
2711 ]);
2712 match parse_stream_event(&error).expect("well-formed stream_error") {
2713 StreamEvent::Error {
2714 stream_id,
2715 code,
2716 message,
2717 } => {
2718 assert_eq!(stream_id, [1u8; 16]);
2719 assert_eq!(code, "cancelled");
2720 assert_eq!(message, "boom");
2721 }
2722 other => panic!("expected Error, got {other:?}"),
2723 }
2724
2725 let reply = Value::Map(vec![
2726 (Value::text("frame_type"), Value::text("stream_reply")),
2727 (Value::text("stream_id"), Value::Bytes(vec![1; 16])),
2728 (Value::text("payload"), Value::text("done")),
2729 (Value::text("responded_by"), Value::Bytes(vec![2; 32])),
2730 ]);
2731 match parse_stream_event(&reply).expect("well-formed stream_reply") {
2732 StreamEvent::Reply {
2733 stream_id,
2734 payload,
2735 responded_by,
2736 } => {
2737 assert_eq!(stream_id, [1u8; 16]);
2738 assert_eq!(payload, Value::text("done"));
2739 assert_eq!(responded_by, [2u8; 32]);
2740 }
2741 other => panic!("expected Reply, got {other:?}"),
2742 }
2743 }
2744
2745 #[test]
2746 fn parse_stream_event_rejects_a_non_stream_frame() {
2747 let frame = Value::Map(vec![
2748 (Value::text("frame_type"), Value::text("call")),
2749 (Value::text("stream_id"), Value::Bytes(vec![1; 16])),
2750 ]);
2751 assert_eq!(
2752 parse_stream_event(&frame).unwrap_err(),
2753 ParseStreamEventError::NotAStreamFrame
2754 );
2755 }
2756
2757 #[test]
2764 fn advertise_frame_matches_the_reference_byte_for_byte() {
2765 let pub_bytes = fixed_array(VECTOR_PUB);
2766 let identity = vector_identity();
2767 let spec = AdvertiseSpec::new(
2768 VECTOR_ZERO_REALM,
2769 "macula_rust_sdk.test_procedure",
2770 pub_bytes,
2771 );
2772 let signed = sign(
2773 advertise_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2774 &identity,
2775 );
2776 let sig = match signed.get("signature") {
2777 Some(Value::Bytes(b)) => b.clone(),
2778 other => panic!("expected a signature field, got {other:?}"),
2779 };
2780 assert_eq!(hex::encode_upper(&sig), "22AE051A542289279A56FB9C8587341232EF48208F9A8641C77F37E1B5D3D26A4B7C30CDCA4AE6E851FEB4E2FBF9C5B2469AFCC7317D59F5D775A05C99E99C0A");
2781 let encoded = encode(&signed).expect("encodable frame");
2782 assert_eq!(encoded.len(), 330);
2783 }
2784
2785 #[test]
2786 fn unadvertise_frame_matches_the_reference_byte_for_byte() {
2787 let pub_bytes = fixed_array(VECTOR_PUB);
2788 let identity = vector_identity();
2789 let spec = UnadvertiseSpec::new(
2790 VECTOR_ZERO_REALM,
2791 "macula_rust_sdk.test_procedure",
2792 pub_bytes,
2793 );
2794 let signed = sign(
2795 unadvertise_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
2796 &identity,
2797 );
2798 let sig = match signed.get("signature") {
2799 Some(Value::Bytes(b)) => b.clone(),
2800 other => panic!("expected a signature field, got {other:?}"),
2801 };
2802 assert_eq!(hex::encode_upper(&sig), "C4111E5C2685DCDDB035B9DA29AD2A30D90BC7CAC09620A675D9A3DB480508FDAD7DCDD145B77607395DBF6195643BBA60C2C6D29E2DCFE5F70F20CF15DA2600");
2803 let encoded = encode(&signed).expect("encodable frame");
2804 assert_eq!(encoded.len(), 323);
2805 }
2806}