use std::time::{SystemTime, UNIX_EPOCH};
use crate::cbor::{self, Value};
use crate::identity::KeyPair;
pub const SIG_DOMAIN: &[u8] = b"macula-v2-frame\0";
pub const PROTOCOL_VERSION: i128 = 2;
pub const MAX_FRAME_BYTES: usize = 0x00FF_FFFF;
fn current_millis() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("system clock is after the Unix epoch")
.as_millis() as u64
}
fn fresh_frame_id() -> [u8; 16] {
*uuid::Uuid::now_v7().as_bytes()
}
fn base(
frame_type: &str,
capabilities: u64,
frame_id: [u8; 16],
sent_at_ms: u64,
) -> Vec<(Value, Value)> {
vec![
(Value::text("version"), Value::Int(PROTOCOL_VERSION)),
(Value::text("frame_type"), Value::text(frame_type)),
(Value::text("frame_id"), Value::Bytes(frame_id.to_vec())),
(Value::text("sent_at_ms"), Value::Int(sent_at_ms as i128)),
(
Value::text("capabilities"),
Value::Int(capabilities as i128),
),
(Value::text("realm"), Value::Null),
(Value::text("call_id"), Value::Null),
(Value::text("source_route"), Value::Null),
]
}
fn bytes32_list(items: &[[u8; 32]]) -> Value {
Value::List(items.iter().map(|b| Value::Bytes(b.to_vec())).collect())
}
#[derive(Debug, Clone)]
pub struct ConnectSpec {
pub node_id: [u8; 32],
pub station_id: [u8; 32],
pub realms: Vec<[u8; 32]>,
pub capabilities: u64,
pub puzzle_evidence: [u8; 32],
pub addresses: Vec<Value>,
pub site: Option<Value>,
pub endorsements: Vec<Value>,
}
impl ConnectSpec {
pub fn new(node_id: [u8; 32], puzzle_evidence: [u8; 32]) -> Self {
Self {
node_id,
station_id: node_id,
realms: Vec::new(),
capabilities: 0,
puzzle_evidence,
addresses: Vec::new(),
site: None,
endorsements: Vec::new(),
}
}
}
fn connect_value(spec: &ConnectSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
let mut fields = base("connect", spec.capabilities, frame_id, sent_at_ms);
fields.push((Value::text("node_id"), Value::Bytes(spec.node_id.to_vec())));
fields.push((
Value::text("station_id"),
Value::Bytes(spec.station_id.to_vec()),
));
fields.push((Value::text("realms"), bytes32_list(&spec.realms)));
fields.push((
Value::text("addresses"),
Value::List(spec.addresses.clone()),
));
fields.push((
Value::text("site"),
spec.site.clone().unwrap_or(Value::Null),
));
fields.push((
Value::text("puzzle_evidence"),
Value::Bytes(spec.puzzle_evidence.to_vec()),
));
fields.push((
Value::text("endorsements"),
Value::List(spec.endorsements.clone()),
));
Value::Map(fields)
}
pub fn connect(spec: &ConnectSpec) -> Value {
connect_value(spec, fresh_frame_id(), current_millis())
}
fn goodbye_value(reason: &str, detail: Option<&str>, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
let mut fields = base("goodbye", 0, frame_id, sent_at_ms);
fields.push((Value::text("reason"), Value::text(reason)));
fields.push((
Value::text("detail"),
detail
.map(|d| Value::Bytes(d.as_bytes().to_vec()))
.unwrap_or(Value::Null),
));
Value::Map(fields)
}
pub fn goodbye(reason: &str, detail: Option<&str>) -> Value {
goodbye_value(reason, detail, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct CallSpec {
pub call_id: [u8; 16],
pub procedure: String,
pub realm: [u8; 32],
pub payload: Value,
pub deadline_ms: i128,
pub caller: [u8; 32],
pub source_route: Vec<u8>,
pub retry_budget: u64,
pub ucan_token: Vec<u8>,
}
impl CallSpec {
pub fn new(
call_id: [u8; 16],
procedure: impl Into<String>,
realm: [u8; 32],
payload: Value,
deadline_ms: i128,
caller: [u8; 32],
) -> Self {
Self {
call_id,
procedure: procedure.into(),
realm,
payload,
deadline_ms,
caller,
source_route: Vec::new(),
retry_budget: 0,
ucan_token: Vec::new(),
}
}
}
fn call_value(spec: &CallSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("call", 0, frame_id, sent_at_ms))
.with_field("realm", Value::Bytes(spec.realm.to_vec()))
.with_field("call_id", Value::Bytes(spec.call_id.to_vec()))
.with_field(
"procedure",
Value::Bytes(spec.procedure.as_bytes().to_vec()),
)
.with_field("payload", spec.payload.clone())
.with_field("deadline_ms", Value::Int(spec.deadline_ms))
.with_field("caller", Value::Bytes(spec.caller.to_vec()))
.with_field("source_route", Value::Bytes(spec.source_route.clone()))
.with_field("retry_budget", Value::Int(spec.retry_budget as i128))
.with_field("ucan_token", Value::Bytes(spec.ucan_token.clone()))
}
pub fn call(spec: &CallSpec) -> Value {
call_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct ResultSpec {
pub call_id: [u8; 16],
pub payload: Value,
pub responded_by: [u8; 32],
pub source_route_reverse: Vec<u8>,
}
impl ResultSpec {
pub fn new(call_id: [u8; 16], payload: Value, responded_by: [u8; 32]) -> Self {
Self {
call_id,
payload,
responded_by,
source_route_reverse: Vec::new(),
}
}
}
fn result_value(spec: &ResultSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("result", 0, frame_id, sent_at_ms))
.with_field("call_id", Value::Bytes(spec.call_id.to_vec()))
.with_field("payload", spec.payload.clone())
.with_field("responded_by", Value::Bytes(spec.responded_by.to_vec()))
.with_field(
"source_route_reverse",
Value::Bytes(spec.source_route_reverse.clone()),
)
}
pub fn result(spec: &ResultSpec) -> Value {
result_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct CallErrorSpec {
pub call_id: [u8; 16],
pub code: crate::bolt4::Code,
pub reported_by: [u8; 32],
pub detail: Option<String>,
pub offending_hop: Option<[u8; 32]>,
pub source_route_partial: Vec<u8>,
}
impl CallErrorSpec {
pub fn new(call_id: [u8; 16], code: crate::bolt4::Code, reported_by: [u8; 32]) -> Self {
Self {
call_id,
code,
reported_by,
detail: None,
offending_hop: None,
source_route_partial: Vec::new(),
}
}
}
fn call_error_value(spec: &CallErrorSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("error", 0, frame_id, sent_at_ms))
.with_field("call_id", Value::Bytes(spec.call_id.to_vec()))
.with_field("code", Value::Int(spec.code.as_u8() as i128))
.with_field("name", Value::text(spec.code.name()))
.with_field("reported_by", Value::Bytes(spec.reported_by.to_vec()))
.with_field(
"detail",
spec.detail
.as_ref()
.map(|d| Value::Bytes(d.as_bytes().to_vec()))
.unwrap_or(Value::Null),
)
.with_field(
"offending_hop",
spec.offending_hop
.map(|h| Value::Bytes(h.to_vec()))
.unwrap_or(Value::Null),
)
.with_field(
"source_route_partial",
Value::Bytes(spec.source_route_partial.clone()),
)
}
pub fn call_error(spec: &CallErrorSpec) -> Value {
call_error_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct CallInfo {
pub call_id: [u8; 16],
pub procedure: String,
pub realm: [u8; 32],
pub payload: Value,
pub deadline_ms: i128,
pub caller: [u8; 32],
pub ucan_token: Vec<u8>,
}
#[derive(Debug, PartialEq, Eq)]
pub enum ParseCallError {
NotACallFrame,
MissingField(&'static str),
WrongFieldType(&'static str),
}
impl std::fmt::Display for ParseCallError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ParseCallError::NotACallFrame => write!(f, "frame_type is not \"call\""),
ParseCallError::MissingField(name) => write!(f, "missing required field {name:?}"),
ParseCallError::WrongFieldType(name) => write!(f, "field {name:?} has the wrong type"),
}
}
}
impl std::error::Error for ParseCallError {}
pub fn parse_call(frame: &Value) -> Result<CallInfo, ParseCallError> {
match frame.get("frame_type") {
Some(Value::Text(t)) if t == "call" => {}
_ => return Err(ParseCallError::NotACallFrame),
}
let call_id = match frame.get("call_id") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseCallError::WrongFieldType("call_id"))?,
Some(_) => return Err(ParseCallError::WrongFieldType("call_id")),
None => return Err(ParseCallError::MissingField("call_id")),
};
let procedure = match frame.get("procedure") {
Some(Value::Bytes(b)) => {
String::from_utf8(b.clone()).map_err(|_| ParseCallError::WrongFieldType("procedure"))?
}
Some(_) => return Err(ParseCallError::WrongFieldType("procedure")),
None => return Err(ParseCallError::MissingField("procedure")),
};
let realm = match frame.get("realm") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseCallError::WrongFieldType("realm"))?,
Some(_) => return Err(ParseCallError::WrongFieldType("realm")),
None => return Err(ParseCallError::MissingField("realm")),
};
let payload = frame
.get("payload")
.cloned()
.ok_or(ParseCallError::MissingField("payload"))?;
let deadline_ms = match frame.get("deadline_ms") {
Some(Value::Int(n)) => *n,
Some(_) => return Err(ParseCallError::WrongFieldType("deadline_ms")),
None => return Err(ParseCallError::MissingField("deadline_ms")),
};
let caller = match frame.get("caller") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseCallError::WrongFieldType("caller"))?,
Some(_) => return Err(ParseCallError::WrongFieldType("caller")),
None => return Err(ParseCallError::MissingField("caller")),
};
let ucan_token = match frame.get("ucan_token") {
Some(Value::Bytes(b)) => b.clone(),
_ => Vec::new(),
};
Ok(CallInfo {
call_id,
procedure,
realm,
payload,
deadline_ms,
caller,
ucan_token,
})
}
#[derive(Debug, Clone)]
pub enum CallResponse {
Result {
payload: Value,
responded_by: [u8; 32],
},
Error {
code: u8,
name: String,
reported_by: [u8; 32],
detail: Option<String>,
},
}
#[derive(Debug, PartialEq, Eq)]
pub enum ParseCallResponseError {
NotAResultOrError,
MissingField(&'static str),
WrongFieldType(&'static str),
}
impl std::fmt::Display for ParseCallResponseError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ParseCallResponseError::NotAResultOrError => {
write!(f, "frame_type is neither \"result\" nor \"error\"")
}
ParseCallResponseError::MissingField(name) => {
write!(f, "missing required field {name:?}")
}
ParseCallResponseError::WrongFieldType(name) => {
write!(f, "field {name:?} has the wrong type")
}
}
}
}
impl std::error::Error for ParseCallResponseError {}
pub fn frame_call_id(frame: &Value) -> Option<[u8; 16]> {
match frame.get("call_id") {
Some(Value::Bytes(b)) => b.as_slice().try_into().ok(),
_ => None,
}
}
pub fn parse_call_response(frame: &Value) -> Result<CallResponse, ParseCallResponseError> {
match frame.get("frame_type") {
Some(Value::Text(t)) if t == "result" => {
let payload = frame
.get("payload")
.cloned()
.ok_or(ParseCallResponseError::MissingField("payload"))?;
let responded_by = get_bytes32_generic(frame, "responded_by")?;
Ok(CallResponse::Result {
payload,
responded_by,
})
}
Some(Value::Text(t)) if t == "error" => {
let code = match frame.get("code") {
Some(Value::Int(n)) if (0..=255).contains(n) => *n as u8,
Some(_) => return Err(ParseCallResponseError::WrongFieldType("code")),
None => return Err(ParseCallResponseError::MissingField("code")),
};
let name = match frame.get("name") {
Some(Value::Text(t)) => t.clone(),
Some(_) => return Err(ParseCallResponseError::WrongFieldType("name")),
None => return Err(ParseCallResponseError::MissingField("name")),
};
let reported_by = get_bytes32_generic(frame, "reported_by")?;
let detail = match frame.get("detail") {
None | Some(Value::Null) => None,
Some(Value::Bytes(b)) => Some(
String::from_utf8(b.clone())
.map_err(|_| ParseCallResponseError::WrongFieldType("detail"))?,
),
Some(_) => return Err(ParseCallResponseError::WrongFieldType("detail")),
};
Ok(CallResponse::Error {
code,
name,
reported_by,
detail,
})
}
_ => Err(ParseCallResponseError::NotAResultOrError),
}
}
fn get_bytes32_generic(
frame: &Value,
field: &'static str,
) -> Result<[u8; 32], ParseCallResponseError> {
match frame.get(field) {
None => Err(ParseCallResponseError::MissingField(field)),
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseCallResponseError::WrongFieldType(field)),
Some(_) => Err(ParseCallResponseError::WrongFieldType(field)),
}
}
#[derive(Debug, Clone)]
pub struct PublishSpec {
pub topic: String,
pub realm: [u8; 32],
pub publisher: [u8; 32],
pub seq: u64,
pub payload: Value,
pub published_at_ms: u64,
pub ttl_ms: Option<u64>,
}
impl PublishSpec {
pub fn new(
topic: impl Into<String>,
realm: [u8; 32],
publisher: [u8; 32],
seq: u64,
payload: Value,
published_at_ms: u64,
) -> Self {
Self {
topic: topic.into(),
realm,
publisher,
seq,
payload,
published_at_ms,
ttl_ms: None,
}
}
}
fn publish_value(spec: &PublishSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("publish", 0, frame_id, sent_at_ms))
.with_field("realm", Value::Bytes(spec.realm.to_vec()))
.with_field("topic", Value::Bytes(spec.topic.as_bytes().to_vec()))
.with_field("publisher", Value::Bytes(spec.publisher.to_vec()))
.with_field("seq", Value::Int(spec.seq as i128))
.with_field("payload", spec.payload.clone())
.with_field("published_at_ms", Value::Int(spec.published_at_ms as i128))
.with_field(
"ttl_ms",
spec.ttl_ms
.map(|t| Value::Int(t as i128))
.unwrap_or(Value::Null),
)
}
pub fn publish(spec: &PublishSpec) -> Value {
publish_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct SubscribeSpec {
pub topic: String,
pub realm: [u8; 32],
pub subscriber: [u8; 32],
}
impl SubscribeSpec {
pub fn new(topic: impl Into<String>, realm: [u8; 32], subscriber: [u8; 32]) -> Self {
Self {
topic: topic.into(),
realm,
subscriber,
}
}
}
fn subscribe_value(spec: &SubscribeSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("subscribe", 0, frame_id, sent_at_ms))
.with_field("realm", Value::Bytes(spec.realm.to_vec()))
.with_field("topic", Value::Bytes(spec.topic.as_bytes().to_vec()))
.with_field("subscriber", Value::Bytes(spec.subscriber.to_vec()))
.with_field("filter", Value::Null)
.with_field("options", Value::Map(vec![]))
}
pub fn subscribe(spec: &SubscribeSpec) -> Value {
subscribe_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct UnsubscribeSpec {
pub topic: String,
pub realm: [u8; 32],
pub subscriber: [u8; 32],
}
impl UnsubscribeSpec {
pub fn new(topic: impl Into<String>, realm: [u8; 32], subscriber: [u8; 32]) -> Self {
Self {
topic: topic.into(),
realm,
subscriber,
}
}
}
fn unsubscribe_value(spec: &UnsubscribeSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("unsubscribe", 0, frame_id, sent_at_ms))
.with_field("realm", Value::Bytes(spec.realm.to_vec()))
.with_field("topic", Value::Bytes(spec.topic.as_bytes().to_vec()))
.with_field("subscriber", Value::Bytes(spec.subscriber.to_vec()))
}
pub fn unsubscribe(spec: &UnsubscribeSpec) -> Value {
unsubscribe_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct EventInfo {
pub topic: String,
pub realm: [u8; 32],
pub publisher: [u8; 32],
pub seq: u64,
pub payload: Value,
pub delivered_via: String,
}
#[derive(Debug, PartialEq, Eq)]
pub enum ParseEventError {
NotAnEventFrame,
MissingField(&'static str),
WrongFieldType(&'static str),
}
impl std::fmt::Display for ParseEventError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ParseEventError::NotAnEventFrame => write!(f, "frame_type is not \"event\""),
ParseEventError::MissingField(name) => write!(f, "missing required field {name:?}"),
ParseEventError::WrongFieldType(name) => write!(f, "field {name:?} has the wrong type"),
}
}
}
impl std::error::Error for ParseEventError {}
pub fn parse_event(frame: &Value) -> Result<EventInfo, ParseEventError> {
match frame.get("frame_type") {
Some(Value::Text(t)) if t == "event" => {}
_ => return Err(ParseEventError::NotAnEventFrame),
}
let topic = match frame.get("topic") {
Some(Value::Bytes(b)) => {
String::from_utf8(b.clone()).map_err(|_| ParseEventError::WrongFieldType("topic"))?
}
Some(_) => return Err(ParseEventError::WrongFieldType("topic")),
None => return Err(ParseEventError::MissingField("topic")),
};
let realm = match frame.get("realm") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseEventError::WrongFieldType("realm"))?,
Some(_) => return Err(ParseEventError::WrongFieldType("realm")),
None => return Err(ParseEventError::MissingField("realm")),
};
let publisher = match frame.get("publisher") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseEventError::WrongFieldType("publisher"))?,
Some(_) => return Err(ParseEventError::WrongFieldType("publisher")),
None => return Err(ParseEventError::MissingField("publisher")),
};
let seq = match frame.get("seq") {
Some(Value::Int(n)) if *n >= 0 => *n as u64,
Some(_) => return Err(ParseEventError::WrongFieldType("seq")),
None => return Err(ParseEventError::MissingField("seq")),
};
let payload = frame
.get("payload")
.cloned()
.ok_or(ParseEventError::MissingField("payload"))?;
let delivered_via = match frame.get("delivered_via") {
Some(Value::Text(t)) => t.clone(),
Some(_) => return Err(ParseEventError::WrongFieldType("delivered_via")),
None => return Err(ParseEventError::MissingField("delivered_via")),
};
Ok(EventInfo {
topic,
realm,
publisher,
seq,
payload,
delivered_via,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct HelloInfo {
pub node_id: [u8; 32],
pub station_id: [u8; 32],
pub realms: Vec<[u8; 32]>,
pub capabilities: u64,
pub accepted: bool,
pub negotiated_capabilities: u64,
pub refusal_code: Option<i128>,
}
#[derive(Debug, PartialEq, Eq)]
pub enum ParseHelloError {
NotAHelloFrame,
MissingField(&'static str),
WrongFieldType(&'static str),
}
impl std::fmt::Display for ParseHelloError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ParseHelloError::NotAHelloFrame => write!(f, "frame_type is not \"hello\""),
ParseHelloError::MissingField(name) => write!(f, "missing required field {name:?}"),
ParseHelloError::WrongFieldType(name) => write!(f, "field {name:?} has the wrong type"),
}
}
}
impl std::error::Error for ParseHelloError {}
fn get_bytes32(frame: &Value, field: &'static str) -> Result<[u8; 32], ParseHelloError> {
match frame.get(field) {
None => Err(ParseHelloError::MissingField(field)),
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseHelloError::WrongFieldType(field)),
Some(_) => Err(ParseHelloError::WrongFieldType(field)),
}
}
fn get_bytes32_list(frame: &Value, field: &'static str) -> Result<Vec<[u8; 32]>, ParseHelloError> {
match frame.get(field) {
None => Err(ParseHelloError::MissingField(field)),
Some(Value::List(items)) => items
.iter()
.map(|v| match v {
Value::Bytes(b) => b
.as_slice()
.try_into()
.map_err(|_| ParseHelloError::WrongFieldType(field)),
_ => Err(ParseHelloError::WrongFieldType(field)),
})
.collect(),
Some(_) => Err(ParseHelloError::WrongFieldType(field)),
}
}
fn get_uint(frame: &Value, field: &'static str) -> Result<u64, ParseHelloError> {
match frame.get(field) {
None => Err(ParseHelloError::MissingField(field)),
Some(Value::Int(n)) if *n >= 0 => Ok(*n as u64),
Some(_) => Err(ParseHelloError::WrongFieldType(field)),
}
}
fn get_bool(frame: &Value, field: &'static str) -> Result<bool, ParseHelloError> {
match frame.get(field) {
None => Err(ParseHelloError::MissingField(field)),
Some(Value::Text(t)) if t == "true" => Ok(true),
Some(Value::Text(t)) if t == "false" => Ok(false),
Some(_) => Err(ParseHelloError::WrongFieldType(field)),
}
}
pub fn parse_hello(frame: &Value) -> Result<HelloInfo, ParseHelloError> {
match frame.get("frame_type") {
Some(Value::Text(t)) if t == "hello" => {}
_ => return Err(ParseHelloError::NotAHelloFrame),
}
let refusal_code = match frame.get("refusal_code") {
None | Some(Value::Null) => None,
Some(Value::Int(n)) => Some(*n),
Some(_) => return Err(ParseHelloError::WrongFieldType("refusal_code")),
};
Ok(HelloInfo {
node_id: get_bytes32(frame, "node_id")?,
station_id: get_bytes32(frame, "station_id")?,
realms: get_bytes32_list(frame, "realms")?,
capabilities: get_uint(frame, "capabilities")?,
accepted: get_bool(frame, "accepted")?,
negotiated_capabilities: get_uint(frame, "negotiated_capabilities")?,
refusal_code,
})
}
pub fn sign(frame: Value, identity: &KeyPair) -> Value {
let signable = signable_bytes(&frame);
let sig = identity.sign(&signable);
frame.with_field("signature", Value::Bytes(sig.to_vec()))
}
fn signable_bytes(frame: &Value) -> Vec<u8> {
let unsigned = frame.without(&["signature", "publisher_sig"]);
let canonical =
cbor::encode(&unsigned).expect("a frame built by this module is always encodable");
let mut out = Vec::with_capacity(SIG_DOMAIN.len() + canonical.len());
out.extend_from_slice(SIG_DOMAIN);
out.extend_from_slice(&canonical);
out
}
#[derive(Debug, PartialEq, Eq)]
pub enum VerifyError {
MissingSignature,
BadSignature,
SignatureInvalid,
}
impl std::fmt::Display for VerifyError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
VerifyError::MissingSignature => write!(f, "frame has no signature field"),
VerifyError::BadSignature => write!(f, "signature field is not 64 bytes"),
VerifyError::SignatureInvalid => write!(f, "signature does not verify against pubkey"),
}
}
}
impl std::error::Error for VerifyError {}
pub fn verify(frame: &Value, pubkey: &[u8; 32]) -> Result<(), VerifyError> {
let sig: [u8; 64] = match frame.get("signature") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| VerifyError::BadSignature)?,
_ => return Err(VerifyError::MissingSignature),
};
let signable = signable_bytes(frame);
if crate::identity::verify(&signable, &sig, pubkey) {
Ok(())
} else {
Err(VerifyError::SignatureInvalid)
}
}
pub const EVENT_PUBLISHER_DOMAIN: &[u8] = b"macula-v2-event-pub\0";
pub fn sign_publisher(frame: Value, identity: &KeyPair) -> Value {
let signable = publisher_signing_bytes(&frame);
let sig = identity.sign(&signable);
frame.with_field("publisher_sig", Value::Bytes(sig.to_vec()))
}
#[derive(Debug, PartialEq, Eq)]
pub enum VerifyPublisherError {
MissingPublisherSig,
BadPublisherSig,
PublisherSigInvalid,
}
impl std::fmt::Display for VerifyPublisherError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
VerifyPublisherError::MissingPublisherSig => {
write!(f, "frame has no publisher_sig field")
}
VerifyPublisherError::BadPublisherSig => {
write!(f, "publisher_sig field is not 64 bytes")
}
VerifyPublisherError::PublisherSigInvalid => write!(
f,
"publisher_sig does not verify against the frame's publisher field"
),
}
}
}
impl std::error::Error for VerifyPublisherError {}
pub fn verify_publisher(frame: &Value) -> Result<(), VerifyPublisherError> {
let sig: [u8; 64] = match frame.get("publisher_sig") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| VerifyPublisherError::BadPublisherSig)?,
_ => return Err(VerifyPublisherError::MissingPublisherSig),
};
let pubkey: [u8; 32] = match frame.get("publisher") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| VerifyPublisherError::BadPublisherSig)?,
_ => return Err(VerifyPublisherError::BadPublisherSig),
};
let signable = publisher_signing_bytes(frame);
if crate::identity::verify(&signable, &sig, &pubkey) {
Ok(())
} else {
Err(VerifyPublisherError::PublisherSigInvalid)
}
}
fn publisher_signing_bytes(frame: &Value) -> Vec<u8> {
let fields = ["topic", "realm", "publisher", "seq", "payload"];
let pairs: Vec<(Value, Value)> = fields
.iter()
.map(|f| {
let v = frame.get(f).cloned().unwrap_or(Value::Null);
(Value::text(*f), v)
})
.collect();
let canonical =
cbor::encode(&Value::Map(pairs)).expect("a frame built by this module is always encodable");
let mut out = Vec::with_capacity(EVENT_PUBLISHER_DOMAIN.len() + canonical.len());
out.extend_from_slice(EVENT_PUBLISHER_DOMAIN);
out.extend_from_slice(&canonical);
out
}
#[derive(Debug, PartialEq, Eq)]
pub enum EncodeFrameError {
TooLarge(usize),
Cbor(cbor::IntOutOfRange),
}
impl std::fmt::Display for EncodeFrameError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
EncodeFrameError::TooLarge(n) => {
write!(
f,
"frame is {n} bytes, exceeding the {MAX_FRAME_BYTES}-byte cap"
)
}
EncodeFrameError::Cbor(e) => write!(f, "{e}"),
}
}
}
impl std::error::Error for EncodeFrameError {}
pub fn encode(frame: &Value) -> Result<Vec<u8>, EncodeFrameError> {
let payload = cbor::encode(frame).map_err(EncodeFrameError::Cbor)?;
if payload.len() > MAX_FRAME_BYTES {
return Err(EncodeFrameError::TooLarge(payload.len()));
}
let mut out = Vec::with_capacity(4 + payload.len());
out.extend_from_slice(&(payload.len() as u32).to_be_bytes());
out.extend_from_slice(&payload);
Ok(out)
}
#[derive(Debug)]
pub enum Decoded {
Frame(Value, usize),
More(usize),
}
#[derive(Debug)]
pub enum DecodeFrameError {
TooLarge(usize),
Cbor(cbor::DecodeError),
}
impl std::fmt::Display for DecodeFrameError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
DecodeFrameError::TooLarge(n) => {
write!(
f,
"claimed frame length {n} exceeds the {MAX_FRAME_BYTES}-byte cap"
)
}
DecodeFrameError::Cbor(e) => write!(f, "{e}"),
}
}
}
impl std::error::Error for DecodeFrameError {}
pub fn decode(buf: &[u8]) -> Result<Decoded, DecodeFrameError> {
if buf.len() < 4 {
return Ok(Decoded::More(4 - buf.len()));
}
let len = u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]) as usize;
if len > MAX_FRAME_BYTES {
return Err(DecodeFrameError::TooLarge(len));
}
if buf.len() < 4 + len {
return Ok(Decoded::More(4 + len - buf.len()));
}
let value = cbor::decode(&buf[4..4 + len]).map_err(DecodeFrameError::Cbor)?;
Ok(Decoded::Frame(value, 4 + len))
}
#[derive(Debug, Clone)]
pub struct AdvertiseSpec {
pub realm: [u8; 32],
pub procedure: String,
pub advertiser: [u8; 32],
}
impl AdvertiseSpec {
pub fn new(realm: [u8; 32], procedure: impl Into<String>, advertiser: [u8; 32]) -> Self {
Self {
realm,
procedure: procedure.into(),
advertiser,
}
}
}
fn advertise_value(spec: &AdvertiseSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("advertise", 0, frame_id, sent_at_ms))
.with_field("realm", Value::Bytes(spec.realm.to_vec()))
.with_field(
"procedure",
Value::Bytes(spec.procedure.as_bytes().to_vec()),
)
.with_field("advertiser", Value::Bytes(spec.advertiser.to_vec()))
.with_field("options", Value::Map(vec![]))
}
pub fn advertise(spec: &AdvertiseSpec) -> Value {
advertise_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct UnadvertiseSpec {
pub realm: [u8; 32],
pub procedure: String,
pub advertiser: [u8; 32],
}
impl UnadvertiseSpec {
pub fn new(realm: [u8; 32], procedure: impl Into<String>, advertiser: [u8; 32]) -> Self {
Self {
realm,
procedure: procedure.into(),
advertiser,
}
}
}
fn unadvertise_value(spec: &UnadvertiseSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("unadvertise", 0, frame_id, sent_at_ms))
.with_field("realm", Value::Bytes(spec.realm.to_vec()))
.with_field(
"procedure",
Value::Bytes(spec.procedure.as_bytes().to_vec()),
)
.with_field("advertiser", Value::Bytes(spec.advertiser.to_vec()))
}
pub fn unadvertise(spec: &UnadvertiseSpec) -> Value {
unadvertise_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StreamMode {
ServerStream,
ClientStream,
Bidi,
}
impl StreamMode {
pub fn name(self) -> &'static str {
match self {
StreamMode::ServerStream => "server_stream",
StreamMode::ClientStream => "client_stream",
StreamMode::Bidi => "bidi",
}
}
fn from_name(name: &str) -> Option<Self> {
match name {
"server_stream" => Some(StreamMode::ServerStream),
"client_stream" => Some(StreamMode::ClientStream),
"bidi" => Some(StreamMode::Bidi),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StreamEncoding {
Raw,
Msgpack,
}
impl StreamEncoding {
pub fn name(self) -> &'static str {
match self {
StreamEncoding::Raw => "raw",
StreamEncoding::Msgpack => "msgpack",
}
}
fn from_name(name: &str) -> Option<Self> {
match name {
"raw" => Some(StreamEncoding::Raw),
"msgpack" => Some(StreamEncoding::Msgpack),
_ => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StreamRole {
Send,
Both,
}
impl StreamRole {
pub fn name(self) -> &'static str {
match self {
StreamRole::Send => "send",
StreamRole::Both => "both",
}
}
fn from_name(name: &str) -> Option<Self> {
match name {
"send" => Some(StreamRole::Send),
"both" => Some(StreamRole::Both),
_ => None,
}
}
}
#[derive(Debug, Clone)]
pub struct StreamOpenSpec {
pub stream_id: [u8; 16],
pub procedure: String,
pub realm: [u8; 32],
pub mode: StreamMode,
pub args: Value,
pub deadline_ms: i128,
pub caller: [u8; 32],
pub source_route: Vec<u8>,
pub retry_budget: u64,
}
impl StreamOpenSpec {
pub fn new(
stream_id: [u8; 16],
procedure: impl Into<String>,
realm: [u8; 32],
mode: StreamMode,
args: Value,
deadline_ms: i128,
caller: [u8; 32],
) -> Self {
Self {
stream_id,
procedure: procedure.into(),
realm,
mode,
args,
deadline_ms,
caller,
source_route: Vec::new(),
retry_budget: 0,
}
}
}
fn stream_open_value(spec: &StreamOpenSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("stream_open", 0, frame_id, sent_at_ms))
.with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
.with_field(
"procedure",
Value::Bytes(spec.procedure.as_bytes().to_vec()),
)
.with_field("realm", Value::Bytes(spec.realm.to_vec()))
.with_field("mode", Value::text(spec.mode.name()))
.with_field("args", spec.args.clone())
.with_field("deadline_ms", Value::Int(spec.deadline_ms))
.with_field("caller", Value::Bytes(spec.caller.to_vec()))
.with_field("source_route", Value::Bytes(spec.source_route.clone()))
.with_field("retry_budget", Value::Int(spec.retry_budget as i128))
}
pub fn stream_open(spec: &StreamOpenSpec) -> Value {
stream_open_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct StreamOpenInfo {
pub stream_id: [u8; 16],
pub procedure: String,
pub realm: [u8; 32],
pub mode: StreamMode,
pub args: Value,
pub deadline_ms: i128,
pub caller: [u8; 32],
}
#[derive(Debug, PartialEq, Eq)]
pub enum ParseStreamOpenError {
NotAStreamOpenFrame,
MissingField(&'static str),
WrongFieldType(&'static str),
}
impl std::fmt::Display for ParseStreamOpenError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ParseStreamOpenError::NotAStreamOpenFrame => {
write!(f, "frame_type is not \"stream_open\"")
}
ParseStreamOpenError::MissingField(name) => {
write!(f, "missing required field {name:?}")
}
ParseStreamOpenError::WrongFieldType(name) => {
write!(f, "field {name:?} has the wrong type")
}
}
}
}
impl std::error::Error for ParseStreamOpenError {}
pub fn parse_stream_open(frame: &Value) -> Result<StreamOpenInfo, ParseStreamOpenError> {
match frame.get("frame_type") {
Some(Value::Text(t)) if t == "stream_open" => {}
_ => return Err(ParseStreamOpenError::NotAStreamOpenFrame),
}
let stream_id = match frame.get("stream_id") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseStreamOpenError::WrongFieldType("stream_id"))?,
Some(_) => return Err(ParseStreamOpenError::WrongFieldType("stream_id")),
None => return Err(ParseStreamOpenError::MissingField("stream_id")),
};
let procedure = match frame.get("procedure") {
Some(Value::Bytes(b)) => String::from_utf8(b.clone())
.map_err(|_| ParseStreamOpenError::WrongFieldType("procedure"))?,
Some(_) => return Err(ParseStreamOpenError::WrongFieldType("procedure")),
None => return Err(ParseStreamOpenError::MissingField("procedure")),
};
let realm = match frame.get("realm") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseStreamOpenError::WrongFieldType("realm"))?,
Some(_) => return Err(ParseStreamOpenError::WrongFieldType("realm")),
None => return Err(ParseStreamOpenError::MissingField("realm")),
};
let mode = match frame.get("mode") {
Some(Value::Text(t)) => {
StreamMode::from_name(t).ok_or(ParseStreamOpenError::WrongFieldType("mode"))?
}
Some(_) => return Err(ParseStreamOpenError::WrongFieldType("mode")),
None => return Err(ParseStreamOpenError::MissingField("mode")),
};
let args = frame
.get("args")
.cloned()
.ok_or(ParseStreamOpenError::MissingField("args"))?;
let deadline_ms = match frame.get("deadline_ms") {
Some(Value::Int(n)) => *n,
Some(_) => return Err(ParseStreamOpenError::WrongFieldType("deadline_ms")),
None => return Err(ParseStreamOpenError::MissingField("deadline_ms")),
};
let caller = match frame.get("caller") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseStreamOpenError::WrongFieldType("caller"))?,
Some(_) => return Err(ParseStreamOpenError::WrongFieldType("caller")),
None => return Err(ParseStreamOpenError::MissingField("caller")),
};
Ok(StreamOpenInfo {
stream_id,
procedure,
realm,
mode,
args,
deadline_ms,
caller,
})
}
#[derive(Debug, Clone)]
pub struct StreamDataSpec {
pub stream_id: [u8; 16],
pub seq: u64,
pub encoding: StreamEncoding,
pub body: Value,
pub signer: Option<[u8; 32]>,
}
impl StreamDataSpec {
pub fn new(
stream_id: [u8; 16],
seq: u64,
encoding: StreamEncoding,
body: Value,
signer: Option<[u8; 32]>,
) -> Self {
Self {
stream_id,
seq,
encoding,
body,
signer,
}
}
}
fn stream_data_value(spec: &StreamDataSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
let value = Value::Map(base("stream_data", 0, frame_id, sent_at_ms))
.with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
.with_field("seq", Value::Int(spec.seq as i128))
.with_field("encoding", Value::text(spec.encoding.name()))
.with_field("body", spec.body.clone());
with_optional_signer(value, spec.signer)
}
pub fn stream_data(spec: &StreamDataSpec) -> Value {
stream_data_value(spec, fresh_frame_id(), current_millis())
}
fn with_optional_signer(value: Value, signer: Option<[u8; 32]>) -> Value {
match signer {
Some(pub_key) => value.with_field("signer", Value::Bytes(pub_key.to_vec())),
None => value,
}
}
#[derive(Debug, Clone)]
pub struct StreamEndSpec {
pub stream_id: [u8; 16],
pub role: StreamRole,
pub signer: Option<[u8; 32]>,
}
impl StreamEndSpec {
pub fn new(stream_id: [u8; 16], role: StreamRole, signer: Option<[u8; 32]>) -> Self {
Self {
stream_id,
role,
signer,
}
}
}
fn stream_end_value(spec: &StreamEndSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
let value = Value::Map(base("stream_end", 0, frame_id, sent_at_ms))
.with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
.with_field("role", Value::text(spec.role.name()));
with_optional_signer(value, spec.signer)
}
pub fn stream_end(spec: &StreamEndSpec) -> Value {
stream_end_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct StreamErrorSpec {
pub stream_id: [u8; 16],
pub code: String,
pub message: String,
pub signer: Option<[u8; 32]>,
}
impl StreamErrorSpec {
pub fn new(
stream_id: [u8; 16],
code: impl Into<String>,
message: impl Into<String>,
signer: Option<[u8; 32]>,
) -> Self {
Self {
stream_id,
code: code.into(),
message: message.into(),
signer,
}
}
}
fn stream_error_value(spec: &StreamErrorSpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
let value = Value::Map(base("stream_error", 0, frame_id, sent_at_ms))
.with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
.with_field("code", Value::Bytes(spec.code.as_bytes().to_vec()))
.with_field("message", Value::Bytes(spec.message.as_bytes().to_vec()));
with_optional_signer(value, spec.signer)
}
pub fn stream_error(spec: &StreamErrorSpec) -> Value {
stream_error_value(spec, fresh_frame_id(), current_millis())
}
#[derive(Debug, Clone)]
pub struct StreamReplySpec {
pub stream_id: [u8; 16],
pub payload: Value,
pub responded_by: [u8; 32],
}
impl StreamReplySpec {
pub fn new(stream_id: [u8; 16], payload: Value, responded_by: [u8; 32]) -> Self {
Self {
stream_id,
payload,
responded_by,
}
}
}
fn stream_reply_value(spec: &StreamReplySpec, frame_id: [u8; 16], sent_at_ms: u64) -> Value {
Value::Map(base("stream_reply", 0, frame_id, sent_at_ms))
.with_field("stream_id", Value::Bytes(spec.stream_id.to_vec()))
.with_field("payload", spec.payload.clone())
.with_field("responded_by", Value::Bytes(spec.responded_by.to_vec()))
}
pub fn stream_reply(spec: &StreamReplySpec) -> Value {
stream_reply_value(spec, fresh_frame_id(), current_millis())
}
pub fn frame_stream_id(frame: &Value) -> Option<[u8; 16]> {
match frame.get("stream_id") {
Some(Value::Bytes(b)) => b.as_slice().try_into().ok(),
_ => None,
}
}
#[derive(Debug, Clone)]
pub enum StreamEvent {
Data {
stream_id: [u8; 16],
seq: u64,
encoding: StreamEncoding,
body: Value,
},
End {
stream_id: [u8; 16],
role: StreamRole,
},
Error {
stream_id: [u8; 16],
code: String,
message: String,
},
Reply {
stream_id: [u8; 16],
payload: Value,
responded_by: [u8; 32],
},
}
#[derive(Debug, PartialEq, Eq)]
pub enum ParseStreamEventError {
NotAStreamFrame,
MissingField(&'static str),
WrongFieldType(&'static str),
}
impl std::fmt::Display for ParseStreamEventError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ParseStreamEventError::NotAStreamFrame => write!(
f,
"frame_type is none of stream_data/stream_end/stream_error/stream_reply"
),
ParseStreamEventError::MissingField(name) => {
write!(f, "missing required field {name:?}")
}
ParseStreamEventError::WrongFieldType(name) => {
write!(f, "field {name:?} has the wrong type")
}
}
}
}
impl std::error::Error for ParseStreamEventError {}
pub fn parse_stream_event(frame: &Value) -> Result<StreamEvent, ParseStreamEventError> {
let stream_id = match frame.get("stream_id") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseStreamEventError::WrongFieldType("stream_id"))?,
Some(_) => return Err(ParseStreamEventError::WrongFieldType("stream_id")),
None => return Err(ParseStreamEventError::MissingField("stream_id")),
};
match frame.get("frame_type") {
Some(Value::Text(t)) if t == "stream_data" => {
let seq = match frame.get("seq") {
Some(Value::Int(n)) if *n >= 0 => *n as u64,
Some(_) => return Err(ParseStreamEventError::WrongFieldType("seq")),
None => return Err(ParseStreamEventError::MissingField("seq")),
};
let encoding = match frame.get("encoding") {
Some(Value::Text(t)) => StreamEncoding::from_name(t)
.ok_or(ParseStreamEventError::WrongFieldType("encoding"))?,
Some(_) => return Err(ParseStreamEventError::WrongFieldType("encoding")),
None => return Err(ParseStreamEventError::MissingField("encoding")),
};
let body = frame
.get("body")
.cloned()
.ok_or(ParseStreamEventError::MissingField("body"))?;
Ok(StreamEvent::Data {
stream_id,
seq,
encoding,
body,
})
}
Some(Value::Text(t)) if t == "stream_end" => {
let role = match frame.get("role") {
Some(Value::Text(t)) => {
StreamRole::from_name(t).ok_or(ParseStreamEventError::WrongFieldType("role"))?
}
Some(_) => return Err(ParseStreamEventError::WrongFieldType("role")),
None => return Err(ParseStreamEventError::MissingField("role")),
};
Ok(StreamEvent::End { stream_id, role })
}
Some(Value::Text(t)) if t == "stream_error" => {
let code = match frame.get("code") {
Some(Value::Bytes(b)) => String::from_utf8(b.clone())
.map_err(|_| ParseStreamEventError::WrongFieldType("code"))?,
Some(_) => return Err(ParseStreamEventError::WrongFieldType("code")),
None => return Err(ParseStreamEventError::MissingField("code")),
};
let message = match frame.get("message") {
Some(Value::Bytes(b)) => String::from_utf8(b.clone())
.map_err(|_| ParseStreamEventError::WrongFieldType("message"))?,
Some(_) => return Err(ParseStreamEventError::WrongFieldType("message")),
None => return Err(ParseStreamEventError::MissingField("message")),
};
Ok(StreamEvent::Error {
stream_id,
code,
message,
})
}
Some(Value::Text(t)) if t == "stream_reply" => {
let payload = frame
.get("payload")
.cloned()
.ok_or(ParseStreamEventError::MissingField("payload"))?;
let responded_by = match frame.get("responded_by") {
Some(Value::Bytes(b)) => b
.as_slice()
.try_into()
.map_err(|_| ParseStreamEventError::WrongFieldType("responded_by"))?,
Some(_) => return Err(ParseStreamEventError::WrongFieldType("responded_by")),
None => return Err(ParseStreamEventError::MissingField("responded_by")),
};
Ok(StreamEvent::Reply {
stream_id,
payload,
responded_by,
})
}
Some(_) | None => Err(ParseStreamEventError::NotAStreamFrame),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn hex_bytes(s: &str) -> Vec<u8> {
::hex::decode(s).expect("valid hex fixture")
}
fn fixed_array(hex_str: &str) -> [u8; 32] {
hex_bytes(hex_str).try_into().expect("32-byte fixture")
}
const VECTOR_PUB: &str = "B966A9812649C3D5542FF54954FE090C43FDA6574FE48A0DD326626CFAD29A83";
const VECTOR_PRIV: &str = "457F45FF5A09E172ED15CB20D6CB26B51AD15ED7308C12D478E8631F9CA03D4F";
const VECTOR_PUZZLE_EVIDENCE: &str =
"09D48C91CB46513ED2580BDCEA87C40DA508D4E50EC3DF2F701AFC55D1C5C0B2";
const VECTOR_FRAME_ID: &str = "0192E8B0F1A47000A1B2C3D4E5F60718";
const VECTOR_SENT_AT_MS: u64 = 1_700_000_000_000;
const VECTOR_SIGNATURE: &str = "CF6959A61A2F4D2046F0124C1DD56A6541265F36A24CB18CA8C45C95031854D6AECE5FB93E2AE7BA6C444A09C7C5DED195B6EB0D1CC8E487CCF6E4F0D903B409";
const VECTOR_ENCODED_LEN: usize = 375;
#[test]
fn connect_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = KeyPair::from_seed_bytes(fixed_array(VECTOR_PRIV));
let puzzle_evidence = fixed_array(VECTOR_PUZZLE_EVIDENCE);
let frame_id: [u8; 16] = hex_bytes(VECTOR_FRAME_ID).try_into().expect("16 bytes");
let spec = ConnectSpec::new(pub_bytes, puzzle_evidence);
let unsigned = connect_value(&spec, frame_id, VECTOR_SENT_AT_MS);
let signed = sign(unsigned, &identity);
let sig_field = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(
hex::encode_upper(&sig_field),
VECTOR_SIGNATURE,
"signature diverged from the reference โ canonical CBOR encoding \
or the signing domain/bytes must differ somewhere"
);
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), VECTOR_ENCODED_LEN);
let decoded = match decode(&encoded).expect("valid frame") {
Decoded::Frame(value, consumed) => {
assert_eq!(consumed, encoded.len());
value
}
Decoded::More(n) => panic!("unexpectedly needed {n} more bytes"),
};
verify(&decoded, &pub_bytes).expect("our own signature must verify");
}
#[test]
fn verify_rejects_a_tampered_field() {
let identity = KeyPair::from_seed_bytes(fixed_array(VECTOR_PRIV));
let pub_bytes = identity.public_bytes();
let spec = ConnectSpec::new(pub_bytes, fixed_array(VECTOR_PUZZLE_EVIDENCE));
let signed = sign(connect(&spec), &identity);
let tampered = signed.with_field("capabilities", Value::Int(999));
assert_eq!(
verify(&tampered, &pub_bytes),
Err(VerifyError::SignatureInvalid)
);
}
#[test]
fn verify_rejects_a_missing_signature() {
let frame = Value::Map(vec![(Value::text("frame_type"), Value::text("connect"))]);
let pubkey = [0u8; 32];
assert_eq!(verify(&frame, &pubkey), Err(VerifyError::MissingSignature));
}
#[test]
fn decode_reports_more_for_a_short_buffer() {
assert!(matches!(decode(&[0, 0]), Ok(Decoded::More(2))));
let mut buf = 10u32.to_be_bytes().to_vec();
buf.extend_from_slice(&[0, 0]);
assert!(matches!(decode(&buf), Ok(Decoded::More(8))));
}
#[test]
fn decode_rejects_a_length_over_the_cap() {
let buf = ((MAX_FRAME_BYTES as u32) + 1).to_be_bytes();
assert!(matches!(
decode(&buf),
Err(DecodeFrameError::TooLarge(n)) if n == MAX_FRAME_BYTES + 1
));
}
#[test]
fn goodbye_frame_round_trips() {
let frame = goodbye("normal", Some("bye"));
assert_eq!(frame.get("frame_type"), Some(&Value::text("goodbye")));
assert_eq!(frame.get("reason"), Some(&Value::text("normal")));
assert_eq!(frame.get("detail"), Some(&Value::Bytes(b"bye".to_vec())));
}
#[test]
fn goodbye_without_detail_is_null() {
let frame = goodbye("timeout", None);
assert_eq!(frame.get("detail"), Some(&Value::Null));
}
#[test]
fn parse_hello_reads_a_well_formed_frame() {
let node_id = [7u8; 32];
let station_id = [8u8; 32];
let realm = [9u8; 32];
let hello = Value::Map(vec![
(Value::text("frame_type"), Value::text("hello")),
(Value::text("node_id"), Value::Bytes(node_id.to_vec())),
(Value::text("station_id"), Value::Bytes(station_id.to_vec())),
(
Value::text("realms"),
Value::List(vec![Value::Bytes(realm.to_vec())]),
),
(Value::text("capabilities"), Value::Int(0)),
(Value::text("accepted"), Value::text("true")),
(Value::text("negotiated_capabilities"), Value::Int(3)),
]);
let info = parse_hello(&hello).expect("well-formed hello");
assert_eq!(info.node_id, node_id);
assert_eq!(info.station_id, station_id);
assert_eq!(info.realms, vec![realm]);
assert!(info.accepted);
assert_eq!(info.negotiated_capabilities, 3);
assert_eq!(info.refusal_code, None);
}
#[test]
fn parse_hello_rejects_the_wrong_frame_type() {
let frame = Value::Map(vec![(Value::text("frame_type"), Value::text("connect"))]);
assert_eq!(parse_hello(&frame), Err(ParseHelloError::NotAHelloFrame));
}
#[test]
fn parse_hello_reports_a_missing_field() {
let frame = Value::Map(vec![(Value::text("frame_type"), Value::text("hello"))]);
assert_eq!(
parse_hello(&frame),
Err(ParseHelloError::MissingField("node_id"))
);
}
const VECTOR_CALL_ID: &str = "AABBCCDDEEFF00112233445566778899";
const VECTOR_ZERO_REALM: [u8; 32] = [0u8; 32];
fn vector_identity() -> KeyPair {
KeyPair::from_seed_bytes(fixed_array(VECTOR_PRIV))
}
fn vector_call_id() -> [u8; 16] {
hex_bytes(VECTOR_CALL_ID).try_into().expect("16 bytes")
}
fn vector_frame_id() -> [u8; 16] {
hex_bytes(VECTOR_FRAME_ID).try_into().expect("16 bytes")
}
#[test]
fn call_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = CallSpec::new(
vector_call_id(),
"_content.get_manifest",
VECTOR_ZERO_REALM,
Value::Map(vec![(Value::text("hello"), Value::text("world"))]),
1_700_000_030_000,
pub_bytes,
);
let signed = sign(
call_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(
hex::encode_upper(&sig),
"A6BC174F0241E644F634702C08781C8FC8BD3CDE3CA9650DE8A731A01203D9B9403A2CAD75800F7B8C9AAE16FA146B1195FF03F0E6DC4595A652D7F29BFE350A"
);
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 386);
}
#[test]
fn result_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = ResultSpec::new(vector_call_id(), Value::text("ok-result"), pub_bytes);
let signed = sign(
result_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(
hex::encode_upper(&sig),
"03E8F72D51D958C318B7F1C25D78408408317DEAB23434D6EA32F211CADEA1C62900DA15AFF603E795B19A388D382BDB10E65AEFC6F0CE551270AB172A88E50B"
);
assert_eq!(encode(&signed).expect("encodable").len(), 301);
}
#[test]
fn error_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = CallErrorSpec::new(
vector_call_id(),
crate::bolt4::Code::UnknownNextPeer,
pub_bytes,
);
let signed = sign(
call_error_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(
hex::encode_upper(&sig),
"182ECD5217CE378F576635B23CC8C9F265555142845D6CBA033A282BAED97966C23FBE91D08507FB8E840375AA17665763804F40F89102F8D3EDAD4DA98FC20D"
);
assert_eq!(encode(&signed).expect("encodable").len(), 333);
}
#[test]
fn publish_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = PublishSpec::new(
"test.topic",
VECTOR_ZERO_REALM,
pub_bytes,
42,
Value::text("published-data"),
VECTOR_SENT_AT_MS,
);
let signed = sign(
publish_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(
hex::encode_upper(&sig),
"DD49D10EFA9F2EED0A393DC02DC5BBAC25D6731562EA39F5AB2E5337824527AFFBC7D917AF4DE5EFDBE5BC41E58659E05EC6FDE4E91FB1A32CC9C211456DF10C"
);
assert_eq!(encode(&signed).expect("encodable").len(), 355);
}
#[test]
fn publisher_sig_matches_the_erlang_reference() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = PublishSpec::new(
"acme/svc.do",
VECTOR_ZERO_REALM,
pub_bytes,
42,
Value::Bytes(b"hello".to_vec()),
VECTOR_SENT_AT_MS,
);
let unsigned = publish_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS);
let with_pub_sig = sign_publisher(unsigned, &identity);
let sig = match with_pub_sig.get("publisher_sig") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a publisher_sig field, got {other:?}"),
};
assert_eq!(
hex::encode_upper(&sig),
"C11BEB676A590FD1BA86F0B77E377B4582AA461DB1283F64E57224E920A7BD0A2C7D36271B795FFC3CB4F2C7BB8925B034431AA6425E25B2AEEFAC026883BB0C"
);
verify_publisher(&with_pub_sig).expect("our own freshly-signed frame must verify");
let tampered = with_pub_sig
.clone()
.with_field("payload", Value::Bytes(b"world".to_vec()));
assert!(
verify_publisher(&tampered).is_err(),
"verify_publisher accepted a frame with a tampered payload"
);
assert_eq!(
verify_publisher(&unsigned_publish_for_tamper_check(&spec)),
Err(VerifyPublisherError::MissingPublisherSig)
);
}
fn unsigned_publish_for_tamper_check(spec: &PublishSpec) -> Value {
publish_value(spec, vector_frame_id(), VECTOR_SENT_AT_MS)
}
#[test]
fn publish_frame_with_both_signatures_round_trips() {
let identity = KeyPair::generate();
let pub_bytes = identity.node_id();
let spec = PublishSpec::new(
"acme/svc.do",
VECTOR_ZERO_REALM,
pub_bytes,
1,
Value::Bytes(b"hello".to_vec()),
VECTOR_SENT_AT_MS,
);
let unsigned = publish(&spec);
let with_pub_sig = sign_publisher(unsigned, &identity);
let fully_signed = sign(with_pub_sig, &identity);
let encoded = encode(&fully_signed).expect("encodable");
let decoded = match decode(&encoded).expect("decodable") {
Decoded::Frame(value, consumed) => {
assert_eq!(consumed, encoded.len());
value
}
Decoded::More(n) => panic!("unexpectedly needed {n} more bytes"),
};
verify(&decoded, &pub_bytes).expect("per-hop verify on decoded frame");
verify_publisher(&decoded).expect("verify_publisher on decoded frame");
assert!(
decoded.get("publisher_sig").is_some(),
"decoded frame lost publisher_sig"
);
assert!(
decoded.get("signature").is_some(),
"decoded frame lost signature"
);
}
#[test]
fn subscribe_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = SubscribeSpec::new("test.topic", VECTOR_ZERO_REALM, pub_bytes);
let signed = sign(
subscribe_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(
hex::encode_upper(&sig),
"ABDD7304B887A53B149CE4D4C62F1AFD20AE07D8612B76F22006FA6676B8DDB37C1D5106358D32080246BA4355A9E04BF49F73600E752F5F9037D7A93A47020A"
);
assert_eq!(encode(&signed).expect("encodable").len(), 313);
}
#[test]
fn unsubscribe_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = UnsubscribeSpec::new("test.topic", VECTOR_ZERO_REALM, pub_bytes);
let signed = sign(
unsubscribe_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(
hex::encode_upper(&sig),
"C917068BE4E1C5A3C753F249037DD8F44293D888BB252BF1E828671969547969982160C91A0E3CA1C31DE29ED39E3677E7F20F4BDE61539D4618B3703018E403"
);
assert_eq!(encode(&signed).expect("encodable").len(), 298);
}
#[test]
fn event_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let fields = base("event", 0, vector_frame_id(), VECTOR_SENT_AT_MS);
let unsigned = Value::Map(fields)
.with_field("realm", Value::Bytes(VECTOR_ZERO_REALM.to_vec()))
.with_field("topic", Value::Bytes(b"test.topic".to_vec()))
.with_field("publisher", Value::Bytes(pub_bytes.to_vec()))
.with_field("seq", Value::Int(42))
.with_field("payload", Value::text("published-data"))
.with_field("delivered_via", Value::text("direct"));
let signed = sign(unsigned, &identity);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(
hex::encode_upper(&sig),
"9B9EE4EAC375FBD0C9B5A5BC6D82E35739F8ECBF594979891BF35E5BDB53A148B3936AF99217C3D8C12E2EEA0686F68D5FE63284BE6B142F87BFF319DDDB780F"
);
assert_eq!(encode(&signed).expect("encodable").len(), 341);
let decoded = decode(&encode(&signed).unwrap()).unwrap();
let Decoded::Frame(value, _) = decoded else {
panic!("expected a complete frame")
};
let info = parse_event(&value).expect("well-formed event");
assert_eq!(info.topic, "test.topic");
assert_eq!(info.seq, 42);
assert_eq!(info.delivered_via, "direct");
}
#[test]
fn parse_call_response_reads_a_result() {
let frame = Value::Map(vec![
(Value::text("frame_type"), Value::text("result")),
(Value::text("call_id"), Value::Bytes(vec![1; 16])),
(Value::text("payload"), Value::text("ok")),
(Value::text("responded_by"), Value::Bytes(vec![2; 32])),
]);
match parse_call_response(&frame).expect("well-formed result") {
CallResponse::Result {
payload,
responded_by,
} => {
assert_eq!(payload, Value::text("ok"));
assert_eq!(responded_by, [2u8; 32]);
}
other => panic!("expected Result, got {other:?}"),
}
}
#[test]
fn parse_call_response_reads_an_error() {
let frame = Value::Map(vec![
(Value::text("frame_type"), Value::text("error")),
(Value::text("call_id"), Value::Bytes(vec![1; 16])),
(Value::text("code"), Value::Int(1)),
(Value::text("name"), Value::text("unknown_next_peer")),
(Value::text("reported_by"), Value::Bytes(vec![2; 32])),
(Value::text("detail"), Value::Null),
]);
match parse_call_response(&frame).expect("well-formed error") {
CallResponse::Error {
code,
name,
reported_by,
detail,
} => {
assert_eq!(code, 1);
assert_eq!(name, "unknown_next_peer");
assert_eq!(reported_by, [2u8; 32]);
assert_eq!(detail, None);
}
other => panic!("expected Error, got {other:?}"),
}
}
#[test]
fn frame_call_id_reads_from_any_frame_type() {
let frame = Value::Map(vec![(Value::text("call_id"), Value::Bytes(vec![9; 16]))]);
assert_eq!(frame_call_id(&frame), Some([9u8; 16]));
let wrong_size = Value::Map(vec![(Value::text("call_id"), Value::Bytes(vec![9; 32]))]);
assert_eq!(frame_call_id(&wrong_size), None);
}
const VECTOR_STREAM_ID: &str = "0102030405060708090A0B0C0D0E0F10";
fn vector_stream_id() -> [u8; 16] {
hex_bytes(VECTOR_STREAM_ID).try_into().expect("16 bytes")
}
#[test]
fn stream_open_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = StreamOpenSpec::new(
vector_stream_id(),
"macula_rust_sdk.test_stream",
VECTOR_ZERO_REALM,
StreamMode::ClientStream,
Value::Map(vec![(Value::text("hello"), Value::text("world"))]),
1_700_000_030_000,
pub_bytes,
);
let signed = sign(
stream_open_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "6070D8AB71F837591AC2C803C04F9E1D3FA01C9310D33C96A90434820C5E50550F9DEA8A764247EB49AF63447C037E192B7892A365C1A4ACB9BC46B98AA5670F");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 415);
}
#[test]
fn parse_stream_open_round_trips_a_well_formed_frame() {
let pub_bytes = fixed_array(VECTOR_PUB);
let spec = StreamOpenSpec::new(
vector_stream_id(),
"macula_rust_sdk.test_stream",
VECTOR_ZERO_REALM,
StreamMode::ClientStream,
Value::Map(vec![(Value::text("hello"), Value::text("world"))]),
1_700_000_030_000,
pub_bytes,
);
let frame = stream_open_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS);
let info = parse_stream_open(&frame).expect("well-formed stream_open");
assert_eq!(info.stream_id, vector_stream_id());
assert_eq!(info.procedure, "macula_rust_sdk.test_stream");
assert_eq!(info.realm, VECTOR_ZERO_REALM);
assert_eq!(info.mode, StreamMode::ClientStream);
assert_eq!(
info.args,
Value::Map(vec![(Value::text("hello"), Value::text("world"))])
);
assert_eq!(info.deadline_ms, 1_700_000_030_000);
assert_eq!(info.caller, pub_bytes);
}
#[test]
fn parse_stream_open_rejects_the_wrong_frame_type() {
let frame = Value::Map(vec![(
Value::text("frame_type"),
Value::text("stream_data"),
)]);
assert_eq!(
parse_stream_open(&frame).unwrap_err(),
ParseStreamOpenError::NotAStreamOpenFrame
);
}
#[test]
fn stream_data_raw_frame_matches_the_reference_byte_for_byte() {
let identity = vector_identity();
let spec = StreamDataSpec::new(
vector_stream_id(),
0,
StreamEncoding::Raw,
Value::Bytes(b"raw chunk bytes".to_vec()),
None,
);
let signed = sign(
stream_data_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "35770744FE5BD01B86DDA01AB4EF855E4E4FE0EDFEDC89FF690728C585C60A5CB035717E3EA9133C4AD833E226F4DB95E9A5AF9AC59E7BACBB8BDF72611F8003");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 269);
}
#[test]
fn stream_data_with_signer_matches_the_reference_byte_for_byte() {
let identity = vector_identity();
let spec = StreamDataSpec::new(
vector_stream_id(),
0,
StreamEncoding::Raw,
Value::Bytes(b"raw chunk bytes".to_vec()),
Some(fixed_array(VECTOR_PUB)),
);
let signed = sign(
stream_data_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "3EA0B6B6DB1549D2EA42AF015A477FCD6D00B11F48F9CC07AF0914CAC18F22B5C12E5EE446811388F207D688960B67D9BEE7B4D998BE02F2B1426B6C4A06D307");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 310);
}
#[test]
fn stream_data_msgpack_frame_matches_the_reference_byte_for_byte() {
let identity = vector_identity();
let spec = StreamDataSpec::new(
vector_stream_id(),
1,
StreamEncoding::Msgpack,
Value::Map(vec![
(Value::text("a"), Value::Int(1)),
(Value::text("greeting"), Value::Bytes(b"hi".to_vec())),
]),
None,
);
let signed = sign(
stream_data_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "99CA90B0C01FD349DBAF317D03872E5F460426789874D79B6FBE37F4AC92C2AD690A00CDB3734F262D5C58C8F3BFD06F8AE892A8B5655274718A283ABA1D4D08");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 273);
}
#[test]
fn stream_end_frame_matches_the_reference_byte_for_byte() {
let identity = vector_identity();
let spec = StreamEndSpec::new(vector_stream_id(), StreamRole::Send, None);
let signed = sign(
stream_end_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "78F2B94BD5AC70901EABB31D8B17C89B58A88942300C6232545899AFB933B2C4B7399BB183A5660671981B6346DA27033C8F93A99E7EBA96F0F689B03D4F940A");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 239);
}
#[test]
fn stream_end_with_signer_matches_the_reference_byte_for_byte() {
let identity = vector_identity();
let spec = StreamEndSpec::new(
vector_stream_id(),
StreamRole::Send,
Some(fixed_array(VECTOR_PUB)),
);
let signed = sign(
stream_end_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "CC316B0A1C1AD4701AD16D8A140ED62D5DEEFD721C1CEB574CC8755C645CA27413EF9C6A6A9C4768564524C412515C14637A9D6BD4CCB8CD1ADD44F2A240C70C");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 280);
}
#[test]
fn stream_error_frame_matches_the_reference_byte_for_byte() {
let identity = vector_identity();
let spec = StreamErrorSpec::new(vector_stream_id(), "cancelled", "boom", None);
let signed = sign(
stream_error_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "119F379518EC17C603ED5466A57D7AE53198A8AC4D5CA9849934A78994428CB3DAD40BC0EFECE1A0C8EEB0ACC28973C0F7E55DE6444827091814AF0715D9FF0B");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 259);
}
#[test]
fn stream_error_with_signer_matches_the_reference_byte_for_byte() {
let identity = vector_identity();
let spec = StreamErrorSpec::new(
vector_stream_id(),
"cancelled",
"boom",
Some(fixed_array(VECTOR_PUB)),
);
let signed = sign(
stream_error_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "223062E2816C5E6DABCF08A0A4FD01F477F2D1D933F2F1FDC971CAB570003DDE8192CC2F8811CE4A2D180B6781AFA64EB4057947E25CF121F745A9654DC23D0A");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 300);
}
#[test]
fn stream_reply_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = StreamReplySpec::new(
vector_stream_id(),
Value::Map(vec![(Value::text("ok"), Value::text("true"))]),
pub_bytes,
);
let signed = sign(
stream_reply_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "ADF57AD58B253F175ADF72E4717E078C62F3E22CBDDBF8DDC0DD8A47CAAA061E8A37C73BAAB91E450D1D8472021B6A0161169D77E9D186C436D3E6580D48C703");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 295);
}
#[test]
fn frame_stream_id_reads_from_any_frame_type() {
let frame = Value::Map(vec![(Value::text("stream_id"), Value::Bytes(vec![9; 16]))]);
assert_eq!(frame_stream_id(&frame), Some([9u8; 16]));
let wrong_size = Value::Map(vec![(Value::text("stream_id"), Value::Bytes(vec![9; 32]))]);
assert_eq!(frame_stream_id(&wrong_size), None);
}
#[test]
fn parse_stream_event_reads_data_end_error_and_reply() {
let data = Value::Map(vec![
(Value::text("frame_type"), Value::text("stream_data")),
(Value::text("stream_id"), Value::Bytes(vec![1; 16])),
(Value::text("seq"), Value::Int(3)),
(Value::text("encoding"), Value::text("raw")),
(Value::text("body"), Value::Bytes(b"hi".to_vec())),
]);
match parse_stream_event(&data).expect("well-formed stream_data") {
StreamEvent::Data {
stream_id,
seq,
encoding,
body,
} => {
assert_eq!(stream_id, [1u8; 16]);
assert_eq!(seq, 3);
assert_eq!(encoding, StreamEncoding::Raw);
assert_eq!(body, Value::Bytes(b"hi".to_vec()));
}
other => panic!("expected Data, got {other:?}"),
}
let end = Value::Map(vec![
(Value::text("frame_type"), Value::text("stream_end")),
(Value::text("stream_id"), Value::Bytes(vec![1; 16])),
(Value::text("role"), Value::text("both")),
]);
match parse_stream_event(&end).expect("well-formed stream_end") {
StreamEvent::End { stream_id, role } => {
assert_eq!(stream_id, [1u8; 16]);
assert_eq!(role, StreamRole::Both);
}
other => panic!("expected End, got {other:?}"),
}
let error = Value::Map(vec![
(Value::text("frame_type"), Value::text("stream_error")),
(Value::text("stream_id"), Value::Bytes(vec![1; 16])),
(Value::text("code"), Value::Bytes(b"cancelled".to_vec())),
(Value::text("message"), Value::Bytes(b"boom".to_vec())),
]);
match parse_stream_event(&error).expect("well-formed stream_error") {
StreamEvent::Error {
stream_id,
code,
message,
} => {
assert_eq!(stream_id, [1u8; 16]);
assert_eq!(code, "cancelled");
assert_eq!(message, "boom");
}
other => panic!("expected Error, got {other:?}"),
}
let reply = Value::Map(vec![
(Value::text("frame_type"), Value::text("stream_reply")),
(Value::text("stream_id"), Value::Bytes(vec![1; 16])),
(Value::text("payload"), Value::text("done")),
(Value::text("responded_by"), Value::Bytes(vec![2; 32])),
]);
match parse_stream_event(&reply).expect("well-formed stream_reply") {
StreamEvent::Reply {
stream_id,
payload,
responded_by,
} => {
assert_eq!(stream_id, [1u8; 16]);
assert_eq!(payload, Value::text("done"));
assert_eq!(responded_by, [2u8; 32]);
}
other => panic!("expected Reply, got {other:?}"),
}
}
#[test]
fn parse_stream_event_rejects_a_non_stream_frame() {
let frame = Value::Map(vec![
(Value::text("frame_type"), Value::text("call")),
(Value::text("stream_id"), Value::Bytes(vec![1; 16])),
]);
assert_eq!(
parse_stream_event(&frame).unwrap_err(),
ParseStreamEventError::NotAStreamFrame
);
}
#[test]
fn advertise_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = AdvertiseSpec::new(
VECTOR_ZERO_REALM,
"macula_rust_sdk.test_procedure",
pub_bytes,
);
let signed = sign(
advertise_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "22AE051A542289279A56FB9C8587341232EF48208F9A8641C77F37E1B5D3D26A4B7C30CDCA4AE6E851FEB4E2FBF9C5B2469AFCC7317D59F5D775A05C99E99C0A");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 330);
}
#[test]
fn unadvertise_frame_matches_the_reference_byte_for_byte() {
let pub_bytes = fixed_array(VECTOR_PUB);
let identity = vector_identity();
let spec = UnadvertiseSpec::new(
VECTOR_ZERO_REALM,
"macula_rust_sdk.test_procedure",
pub_bytes,
);
let signed = sign(
unadvertise_value(&spec, vector_frame_id(), VECTOR_SENT_AT_MS),
&identity,
);
let sig = match signed.get("signature") {
Some(Value::Bytes(b)) => b.clone(),
other => panic!("expected a signature field, got {other:?}"),
};
assert_eq!(hex::encode_upper(&sig), "C4111E5C2685DCDDB035B9DA29AD2A30D90BC7CAC09620A675D9A3DB480508FDAD7DCDD145B77607395DBF6195643BBA60C2C6D29E2DCFE5F70F20CF15DA2600");
let encoded = encode(&signed).expect("encodable frame");
assert_eq!(encoded.len(), 323);
}
}