///|
/// Decoder state for a bounded length-delimited CBOR stream.
pub(all) struct CborFrameDecoder {
pending : Bytes
max_frame_size : Int
frames_seen : Int
bytes_seen : Int
} derive(Debug)
///|
/// Counters returned by stream consumers and useful for telemetry.
pub(all) struct CborFrameStats {
frames : Int
payload_bytes : Int
wire_bytes : Int
} derive(Eq, Debug)
///|
/// Create a decoder that rejects frames larger than the configured limit.
pub fn CborFrameDecoder::new(max_frame_size : Int) -> CborFrameDecoder {
{
pending: Bytes::from_array([]),
max_frame_size: if max_frame_size < 1 {
1
} else {
max_frame_size
},
frames_seen: 0,
bytes_seen: 0,
}
}
///|
/// Return the maximum payload accepted by a decoder.
pub fn cbor_frame_decoder_limit(decoder : CborFrameDecoder) -> Int {
decoder.max_frame_size
}
///|
/// Return the number of complete frames emitted by a decoder.
pub fn cbor_frame_decoder_frames(decoder : CborFrameDecoder) -> Int {
decoder.frames_seen
}
///|
/// Return the number of bytes currently buffered but not yet decoded.
pub fn cbor_frame_decoder_pending_bytes(decoder : CborFrameDecoder) -> Int {
decoder.pending.length()
}
///|
fn append_bytes(lhs : Bytes, rhs : Bytes) -> Bytes {
let buffer = Buffer()
buffer.write_bytes(lhs)
buffer.write_bytes(rhs)
buffer.to_bytes()
}
///|
fn read_frame_length(bytes : Bytes) -> Int {
let value = (bytes[0].to_int() << 24) |
(bytes[1].to_int() << 16) |
(bytes[2].to_int() << 8) |
bytes[3].to_int()
value
}
///|
fn write_frame_length(buffer : Buffer, length : Int) -> Unit {
buffer.write_byte((length >> 24).to_byte())
buffer.write_byte(((length >> 16) & 0xFF).to_byte())
buffer.write_byte(((length >> 8) & 0xFF).to_byte())
buffer.write_byte((length & 0xFF).to_byte())
}
///|
fn frame_error(message : String) -> CborError {
SemanticError("frame: " + message)
}
///|
/// Encode one CBOR value as a length-delimited frame.
pub fn encode_cbor_frame(
value : CborValue,
max_frame_size : Int,
) -> Bytes raise CborError {
let payload = encode(value)
if max_frame_size < 1 {
raise frame_error("maximum frame size must be positive")
}
if payload.length() > max_frame_size {
raise frame_error("payload exceeds maximum frame size")
}
if payload.length() > 0x7FFFFFFF {
raise frame_error("payload is too large for the host platform")
}
let buffer = Buffer()
write_frame_length(buffer, payload.length())
buffer.write_bytes(payload)
buffer.to_bytes()
}
///|
/// Encode a batch of values into a single transport buffer.
pub fn encode_cbor_frames(
values : Array[CborValue],
max_frame_size : Int,
) -> Bytes raise CborError {
let buffer = Buffer()
for value in values {
let frame = encode_cbor_frame(value, max_frame_size)
buffer.write_bytes(frame)
}
buffer.to_bytes()
}
///|
/// Encode frames incrementally into a caller-owned buffer.
pub fn append_cbor_frame(
buffer : Buffer,
value : CborValue,
max_frame_size : Int,
) -> Unit raise CborError {
buffer.write_bytes(encode_cbor_frame(value, max_frame_size))
}
///|
fn frame_decoder_next(
decoder : CborFrameDecoder,
pending : Bytes,
frames : Array[CborValue],
) -> (CborFrameDecoder, Bytes, Array[CborValue]) raise CborError {
let mut rest = pending
let result = frames
let mut completed = decoder.frames_seen
let mut consumed = decoder.bytes_seen
while rest.length() >= 4 {
let length = read_frame_length(rest)
if length > decoder.max_frame_size {
raise frame_error("payload exceeds decoder limit")
}
if length < 1 {
raise frame_error("zero-length payload is not a valid CBOR item")
}
if rest.length() < 4 + length {
break
}
let payload = rest[4:4 + length].to_owned()
result.push(decode(payload))
completed = completed + 1
consumed = consumed + 4 + length
rest = rest[4 + length:].to_owned()
}
let next = {
pending: rest,
max_frame_size: decoder.max_frame_size,
frames_seen: completed,
bytes_seen: consumed,
}
(next, rest, result)
}
///|
/// Feed a chunk and return all complete values plus the updated decoder.
pub fn cbor_frame_decoder_feed(
decoder : CborFrameDecoder,
chunk : Bytes,
) -> (CborFrameDecoder, Array[CborValue]) raise CborError {
let pending = append_bytes(decoder.pending, chunk)
let (next, _, frames) = frame_decoder_next(decoder, pending, [])
(next, frames)
}
///|
/// Signal end-of-input and reject an incomplete frame header or payload.
pub fn cbor_frame_decoder_finish(
decoder : CborFrameDecoder,
) -> Unit raise CborError {
if decoder.pending.length() != 0 {
raise frame_error("stream ended with an incomplete frame")
}
}
///|
/// Decode a complete transport buffer.
pub fn decode_cbor_frames(
bytes : Bytes,
max_frame_size : Int,
) -> Array[CborValue] raise CborError {
let decoder = CborFrameDecoder::new(max_frame_size)
let (finished, values) = cbor_frame_decoder_feed(decoder, bytes)
cbor_frame_decoder_finish(finished)
values
}
///|
/// Decode frames while retaining a stable stream statistics snapshot.
pub fn cbor_frame_decoder_stats(decoder : CborFrameDecoder) -> CborFrameStats {
{
frames: decoder.frames_seen,
payload_bytes: decoder.bytes_seen - decoder.frames_seen * 4,
wire_bytes: decoder.bytes_seen,
}
}
///|
/// Return the wire size of a frame for admission control.
pub fn cbor_frame_wire_size(value : CborValue) -> Int {
encode(value).length() + 4
}
///|
/// Return whether a value can be emitted under a frame limit.
pub fn cbor_value_fits_frame(value : CborValue, max_frame_size : Int) -> Bool {
if max_frame_size < 1 {
false
} else {
encode(value).length() <= max_frame_size
}
}
///|
/// Split a batch into transport-sized frame batches without losing order.
pub fn cbor_partition_frames(
values : Array[CborValue],
max_wire_bytes : Int,
max_frame_size : Int,
) -> Array[Array[CborValue]] raise CborError {
if max_wire_bytes < 5 {
raise frame_error("maximum wire batch size is too small")
}
let batches = []
let mut current = []
let mut current_bytes = 0
for value in values {
let size = cbor_frame_wire_size(value)
if size > max_wire_bytes {
raise frame_error("a frame exceeds the maximum wire batch size")
}
if !cbor_value_fits_frame(value, max_frame_size) {
raise frame_error("a value exceeds the maximum frame size")
}
if current.length() > 0 && current_bytes + size > max_wire_bytes {
batches.push(current)
current = []
current_bytes = 0
}
current.push(value)
current_bytes = current_bytes + size
}
if current.length() > 0 {
batches.push(current)
}
batches
}
///|
/// Count complete frames in a wire buffer without decoding payloads.
pub fn count_cbor_frames(
bytes : Bytes,
max_frame_size : Int,
) -> Int raise CborError {
let mut offset = 0
let mut count = 0
while offset < bytes.length() {
if bytes.length() - offset < 4 {
raise frame_error("stream ended in a partial header")
}
let length = read_frame_length(bytes[offset:].to_owned())
if length < 1 || length > max_frame_size {
raise frame_error("frame length is outside the configured limit")
}
if bytes.length() - offset < 4 + length {
raise frame_error("stream ended in a partial payload")
}
offset = offset + 4 + length
count = count + 1
}
count
}
///|
/// Calculate transport statistics from a complete frame stream.
pub fn inspect_cbor_frames(
bytes : Bytes,
max_frame_size : Int,
) -> CborFrameStats raise CborError {
let values = decode_cbor_frames(bytes, max_frame_size)
let mut payload_bytes = 0
for value in values {
payload_bytes = payload_bytes + encode(value).length()
}
{ frames: values.length(), payload_bytes, wire_bytes: bytes.length() }
}
///|
/// Add a frame to a buffer and return the resulting statistics.
pub fn append_frame_with_stats(
buffer : Buffer,
value : CborValue,
max_frame_size : Int,
stats : CborFrameStats,
) -> CborFrameStats raise CborError {
append_cbor_frame(buffer, value, max_frame_size)
let payload_size = encode(value).length()
{
frames: stats.frames + 1,
payload_bytes: stats.payload_bytes + payload_size,
wire_bytes: stats.wire_bytes + payload_size + 4,
}
}