///|
/// An ordered group of application messages for queue and log transports.
pub(all) struct AppMessageBatch {
producer : String
sequence : UInt64
messages : Array[AppMessage]
} derive(Eq, Debug)
///|
/// Construct a batch with producer and monotonically increasing sequence metadata.
pub fn app_batch_new(
producer~ : String,
sequence~ : UInt64,
messages~ : Array[AppMessage],
) -> AppMessageBatch {
{ producer, sequence, messages }
}
///|
/// Return a copy with one message appended.
pub fn app_batch_append(
batch : AppMessageBatch,
message : AppMessage,
) -> AppMessageBatch {
let messages = []
for item in batch.messages {
messages.push(item)
}
messages.push(message)
{ producer: batch.producer, sequence: batch.sequence, messages }
}
///|
/// Return the next sequence number for a producer.
pub fn app_batch_next_sequence(batch : AppMessageBatch) -> UInt64 {
batch.sequence + 1UL
}
///|
/// Validate batch metadata and every contained envelope.
pub fn validate_app_batch(batch : AppMessageBatch) -> Unit raise CborError {
if batch.producer.length() == 0 {
raise SemanticError("app batch producer cannot be empty")
}
if batch.producer.length() > 256 {
raise SemanticError("app batch producer exceeds size limit")
}
if batch.messages.length() > 1024 {
raise SemanticError("app batch contains too many messages")
}
for message in batch.messages {
validate_app_message(message)
}
}
///|
/// Encode a batch using an explicit versioned map representation.
pub fn app_batch_to_cbor(batch : AppMessageBatch) -> CborValue {
let messages = []
for message in batch.messages {
messages.push(app_message_to_cbor(message))
}
Map([
(Text("version"), Integer(1L)),
(Text("producer"), Text(batch.producer)),
(Text("sequence"), Unsigned(batch.sequence)),
(Text("messages"), Array(messages)),
])
}
///|
fn batch_field(value : CborValue, name : String) -> CborValue raise CborError {
match cbor_object_get(value, name) {
Some(found) => found
None => raise SemanticError("app batch is missing field: " + name)
}
}
///|
/// Decode and validate a batch from a CBOR value.
pub fn app_batch_from_cbor(
value : CborValue,
) -> AppMessageBatch raise CborError {
match value {
Map(_) => {
let version = app_int(batch_field(value, "version"), "version")
if version != 1L {
raise SemanticError("app batch version is unsupported")
}
let producer = app_text(batch_field(value, "producer"), "producer")
let sequence = match batch_field(value, "sequence") {
Unsigned(number) => number
Integer(number) if number >= 0L => number.reinterpret_as_uint64()
_ => raise SemanticError("app batch sequence must be unsigned")
}
let messages = match batch_field(value, "messages") {
Array(values) => {
let result = []
for item in values {
result.push(app_message_from_cbor(item))
}
result
}
_ => raise SemanticError("app batch messages must be an array")
}
let batch = { producer, sequence, messages }
validate_app_batch(batch)
batch
}
_ => raise SemanticError("app batch must be a CBOR map")
}
}
///|
/// Encode a complete application message batch.
pub fn encode_app_batch(batch : AppMessageBatch) -> Bytes raise CborError {
validate_app_batch(batch)
encode(app_batch_to_cbor(batch))
}
///|
/// Decode a complete application message batch.
pub fn decode_app_batch(bytes : Bytes) -> AppMessageBatch raise CborError {
app_batch_from_cbor(decode(bytes))
}
///|
/// Count messages of one lifecycle kind.
pub fn app_batch_kind_count(
batch : AppMessageBatch,
kind : AppMessageKind,
) -> Int {
let mut count = 0
for message in batch.messages {
if message.kind == kind {
count = count + 1
}
}
count
}
///|
/// Select messages of one kind while preserving their order.
pub fn app_batch_filter_kind(
batch : AppMessageBatch,
kind : AppMessageKind,
) -> AppMessageBatch {
let messages = []
for message in batch.messages {
if message.kind == kind {
messages.push(message)
}
}
{ producer: batch.producer, sequence: batch.sequence, messages }
}
///|
/// Find a message in a batch by its stable ID.
pub fn app_batch_find(batch : AppMessageBatch, id : String) -> AppMessage? {
for message in batch.messages {
if message.id == id {
return Some(message)
}
}
None
}
///|
/// Return all message IDs for audit and deduplication.
pub fn app_batch_ids(batch : AppMessageBatch) -> Array[String] {
let ids = []
for message in batch.messages {
ids.push(message.id)
}
ids
}
///|
/// Return true when all message IDs in the batch are unique.
pub fn app_batch_has_unique_ids(batch : AppMessageBatch) -> Bool {
for i = 0; i < batch.messages.length(); i = i + 1 {
for j = i + 1; j < batch.messages.length(); j = j + 1 {
if batch.messages[i].id == batch.messages[j].id {
return false
}
}
}
true
}
///|
/// Return the encoded size of a batch for transport admission control.
pub fn app_batch_encoded_size(batch : AppMessageBatch) -> Int {
encode(app_batch_to_cbor(batch)).length()
}
///|
/// Return a compact audit summary for a batch.
pub fn app_batch_summary(batch : AppMessageBatch) -> String {
"producer=\{batch.producer} sequence=\{batch.sequence} messages=\{batch.messages.length()}"
}
///|
/// Validate an upper bound and ID uniqueness in one operation.
pub fn validate_app_batch_for_transport(
batch : AppMessageBatch,
max_encoded_bytes : Int,
) -> Unit raise CborError {
validate_app_batch(batch)
if !app_batch_has_unique_ids(batch) {
raise SemanticError("app batch contains duplicate message IDs")
}
if app_batch_encoded_size(batch) > max_encoded_bytes {
raise SemanticError("app batch exceeds transport size limit")
}
}
///|
/// Return whether a batch contains no messages.
pub fn app_batch_is_empty(batch : AppMessageBatch) -> Bool {
batch.messages.length() == 0
}
///|
/// Return the first message in a batch, if any.
pub fn app_batch_first(batch : AppMessageBatch) -> AppMessage? {
if batch.messages.length() == 0 {
None
} else {
Some(batch.messages[0])
}
}
///|
/// Return the last message in a batch, if any.
pub fn app_batch_last(batch : AppMessageBatch) -> AppMessage? {
if batch.messages.length() == 0 {
None
} else {
Some(batch.messages[batch.messages.length() - 1])
}
}
///|
/// Count all CBOR nodes carried by batch payloads.
pub fn app_batch_payload_nodes(batch : AppMessageBatch) -> Int {
let mut count = 0
for message in batch.messages {
count = count + cbor_value_node_count(message.payload)
}
count
}
///|
/// Return the total encoded payload bytes without envelope overhead.
pub fn app_batch_payload_bytes(batch : AppMessageBatch) -> Int {
let mut total = 0
for message in batch.messages {
total = total + encode(message.payload).length()
}
total
}
///|
/// Return all response messages correlated to a request ID.
pub fn app_batch_responses_for(
batch : AppMessageBatch,
request_id : String,
) -> Array[AppMessage] {
let responses = []
for message in batch.messages {
match message.correlation_id {
Some(id) =>
if id == request_id && message.kind == Response {
responses.push(message)
}
None => ()
}
}
responses
}
///|
/// Return the number of request/response pairs represented in a batch.
pub fn app_batch_correlated_response_count(batch : AppMessageBatch) -> Int {
let mut count = 0
for message in batch.messages {
if message.kind == Response && message.correlation_id is Some(_) {
count = count + 1
}
}
count
}
///|
/// Return a deterministic kind histogram for operational dashboards.
pub fn app_batch_kind_summary(batch : AppMessageBatch) -> String {
"request=\{app_batch_kind_count(batch, Request)} response=\{app_batch_kind_count(batch, Response)} event=\{app_batch_kind_count(batch, Event)} command=\{app_batch_kind_count(batch, Command)} notification=\{app_batch_kind_count(batch, Notification)}"
}
///|
/// Return true when every payload remains under a node budget.
pub fn app_batch_payloads_within_budget(
batch : AppMessageBatch,
max_nodes : Int,
) -> Bool {
for message in batch.messages {
if cbor_value_node_count(message.payload) > max_nodes {
return false
}
}
true
}