///|
/// Application-level message kinds for services built on top of CBOR.
pub(all) enum AppMessageKind {
Request
Response
Event
Command
Notification
} derive(Eq, Debug)
///|
/// A versioned message envelope suitable for RPC, events, and queue records.
pub(all) struct AppMessage {
version : Int
kind : AppMessageKind
id : String
correlation_id : String?
headers : Array[(String, String)]
payload : CborValue
} derive(Eq, Debug)
///|
/// Create a message with no correlation ID or headers.
pub fn AppMessage::new(
version~ : Int,
kind~ : AppMessageKind,
id~ : String,
payload~ : CborValue,
) -> AppMessage {
{ version, kind, id, correlation_id: None, headers: [], payload }
}
///|
/// Return a copy with a correlation ID used to link a response to a request.
pub fn app_message_with_correlation(
message : AppMessage,
correlation_id : String,
) -> AppMessage {
{
version: message.version,
kind: message.kind,
id: message.id,
correlation_id: Some(correlation_id),
headers: message.headers,
payload: message.payload,
}
}
///|
/// Return a copy with one application header appended.
pub fn app_message_with_header(
message : AppMessage,
name : String,
value : String,
) -> AppMessage {
let headers = []
for header in message.headers {
headers.push(header)
}
headers.push((name, value))
{
version: message.version,
kind: message.kind,
id: message.id,
correlation_id: message.correlation_id,
headers,
payload: message.payload,
}
}
///|
/// Find the first header with a given name.
pub fn app_message_header(message : AppMessage, name : String) -> String? {
for header in message.headers {
if header.0 == name {
return Some(header.1)
}
}
None
}
///|
fn message_kind_to_int(kind : AppMessageKind) -> Int {
match kind {
Request => 0
Response => 1
Event => 2
Command => 3
Notification => 4
}
}
///|
fn message_kind_from_int(value : Int64) -> Result[AppMessageKind, CborError] {
match value {
0L => Ok(Request)
1L => Ok(Response)
2L => Ok(Event)
3L => Ok(Command)
4L => Ok(Notification)
_ => Err(SemanticError("app message kind is unknown"))
}
}
///|
fn app_field(
fields : Array[(CborValue, CborValue)],
name : String,
) -> CborValue? {
for field in fields {
match field.0 {
Text(key) => if key == name { return Some(field.1) }
_ => ()
}
}
None
}
///|
fn app_required_field(
fields : Array[(CborValue, CborValue)],
name : String,
) -> CborValue raise CborError {
match app_field(fields, name) {
Some(value) => value
None => raise SemanticError("app message is missing field: " + name)
}
}
///|
fn app_text(value : CborValue, field : String) -> String raise CborError {
match value {
Text(text) => text
_ => raise SemanticError("app message field is not text: " + field)
}
}
///|
fn app_int(value : CborValue, field : String) -> Int64 raise CborError {
match value {
Integer(number) => number
Unsigned(number) =>
if number > 0x7FFFFFFFFFFFFFFFUL {
raise SemanticError("app message integer is out of range: " + field)
} else {
number.reinterpret_as_int64()
}
_ => raise SemanticError("app message field is not an integer: " + field)
}
}
///|
fn app_headers_to_cbor(headers : Array[(String, String)]) -> CborValue {
let values = []
for header in headers {
values.push(
CborValue::Array([CborValue::Text(header.0), CborValue::Text(header.1)]),
)
}
CborValue::Array(values)
}
///|
fn app_headers_from_cbor(
value : CborValue,
) -> Array[(String, String)] raise CborError {
match value {
Array(values) => {
let headers = []
for item in values {
match item {
Array(pair) => {
if pair.length() != 2 {
raise SemanticError("app header must contain two values")
}
let name = app_text(pair[0], "header.name")
let header_value = app_text(pair[1], "header.value")
if name.length() == 0 {
raise SemanticError("app header name cannot be empty")
}
headers.push((name, header_value))
}
_ => raise SemanticError("app headers must be pairs")
}
}
headers
}
_ => raise SemanticError("app message headers must be an array")
}
}
///|
/// Convert an application message into its stable CBOR representation.
pub fn app_message_to_cbor(message : AppMessage) -> CborValue {
let fields = [
(CborValue::Text("version"), CborValue::Integer(message.version.to_int64())),
(
CborValue::Text("kind"),
CborValue::Integer(message_kind_to_int(message.kind).to_int64()),
),
(CborValue::Text("id"), CborValue::Text(message.id)),
(CborValue::Text("headers"), app_headers_to_cbor(message.headers)),
(CborValue::Text("payload"), message.payload),
]
match message.correlation_id {
Some(correlation_id) =>
fields.push(
(CborValue::Text("correlation_id"), CborValue::Text(correlation_id)),
)
None => ()
}
CborValue::Map(fields)
}
///|
/// Validate message invariants without encoding the payload.
pub fn validate_app_message(message : AppMessage) -> Unit raise CborError {
if message.version < 1 {
raise SemanticError("app message version must be positive")
}
if message.id.length() == 0 {
raise SemanticError("app message id cannot be empty")
}
if message.headers.length() > 64 {
raise SemanticError("app message has too many headers")
}
for header in message.headers {
if header.0.length() == 0 {
raise SemanticError("app message header name cannot be empty")
}
if header.0.length() > 128 || header.1.length() > 4096 {
raise SemanticError("app message header exceeds size limit")
}
}
}
///|
/// Decode and validate an application message envelope.
pub fn app_message_from_cbor(value : CborValue) -> AppMessage raise CborError {
match value {
Map(fields) => {
let version = app_int(app_required_field(fields, "version"), "version")
if version > 0x7FFFFFFFFFFFFFFFL || version < 1L {
raise SemanticError("app message version is invalid")
}
let kind_value = app_int(app_required_field(fields, "kind"), "kind")
let kind = match message_kind_from_int(kind_value) {
Ok(value) => value
Err(_) => raise SemanticError("app message kind is invalid")
}
let id = app_text(app_required_field(fields, "id"), "id")
let headers = app_headers_from_cbor(app_required_field(fields, "headers"))
let correlation_id = match app_field(fields, "correlation_id") {
Some(value) => Some(app_text(value, "correlation_id"))
None => None
}
let message = {
version: version.to_int(),
kind,
id,
correlation_id,
headers,
payload: app_required_field(fields, "payload"),
}
validate_app_message(message)
message
}
_ => raise SemanticError("app message must be a CBOR map")
}
}
///|
/// Encode a validated application message.
pub fn encode_app_message(message : AppMessage) -> Bytes raise CborError {
validate_app_message(message)
encode(app_message_to_cbor(message))
}
///|
/// Decode a complete application message.
pub fn decode_app_message(bytes : Bytes) -> AppMessage raise CborError {
app_message_from_cbor(decode(bytes))
}
///|
/// Build a request message with an operation header.
pub fn app_request(
id : String,
operation : String,
payload : CborValue,
) -> AppMessage {
app_message_with_header(
AppMessage::new(version=1, kind=AppMessageKind::Request, id~, payload~),
"operation",
operation,
)
}
///|
/// Build a response correlated to a request.
pub fn app_response(
id : String,
correlation_id : String,
payload : CborValue,
) -> AppMessage {
app_message_with_correlation(
AppMessage::new(version=1, kind=AppMessageKind::Response, id~, payload~),
correlation_id,
)
}
///|
/// Build an event message with a topic header.
pub fn app_event(
id : String,
topic : String,
payload : CborValue,
) -> AppMessage {
app_message_with_header(
AppMessage::new(version=1, kind=AppMessageKind::Event, id~, payload~),
"topic",
topic,
)
}
///|
/// Copy a message while replacing its payload.
pub fn app_message_with_payload(
message : AppMessage,
payload : CborValue,
) -> AppMessage {
{
version: message.version,
kind: message.kind,
id: message.id,
correlation_id: message.correlation_id,
headers: message.headers,
payload,
}
}
///|
/// Return true when two messages can be paired as request and response.
pub fn app_messages_are_correlated(
request : AppMessage,
response : AppMessage,
) -> Bool {
match response.correlation_id {
Some(id) => id == request.id && request.kind == AppMessageKind::Request
None => false
}
}
///|
/// Return a compact application-message summary for logs.
pub fn app_message_summary(message : AppMessage) -> String {
let kind_name = match message.kind {
Request => "request"
Response => "response"
Event => "event"
Command => "command"
Notification => "notification"
}
"\{kind_name}:\{message.id}:v\{message.version}:headers=\{message.headers.length()}"
}