///|
/// Encode a raw DEFLATE payload as required by the Avro OCF `deflate` codec.
fn compress_block(
codec : CompressionCodec,
input : Bytes,
) -> Bytes raise OcfError {
match codec {
Null => input
Deflate => {
let compressed = @yazi.compress(input, @yazi.Raw, @yazi.Default) catch {
_ =>
raise Compression(
codec="deflate",
message="DEFLATE compression failed",
)
}
Bytes::from_array(compressed)
}
}
}
///|
/// Decode a raw DEFLATE payload as required by the Avro OCF `deflate` codec.
fn decompress_block(
codec : CompressionCodec,
input : Bytes,
limits : OcfLimits,
) -> Bytes raise OcfError {
match codec {
Null => input
Deflate => {
let (decompressed, _) = @yazi.decompress(input, @yazi.Raw) catch {
_ =>
raise Compression(
codec="deflate",
message="invalid raw DEFLATE payload",
)
}
if decompressed.length() > limits.max_block_bytes {
raise BlockTooLarge(
length=decompressed.length(),
limit=limits.max_block_bytes,
)
}
Bytes::from_array(decompressed)
}
}
}
///|
/// Build the required OCF metadata. Reserved keys cannot be supplied by an
/// application because that would make the file header ambiguous.
fn header_metadata(header : OcfHeader) -> Map[String, Bytes] raise OcfError {
let metadata = header.metadata()
if metadata.contains("avro.schema") || metadata.contains("avro.codec") {
raise InvalidMetadataKey(
"avro.schema and avro.codec are reserved OCF metadata keys",
)
}
for key, _ in metadata {
if key.is_empty() {
raise InvalidMetadataKey(key)
}
}
metadata["avro.schema"] = @utf8.encode(header.schema().source())
metadata["avro.codec"] = @utf8.encode(header.codec().name())
metadata
}
///|
/// Write the Avro map representation used in an OCF header.
fn write_metadata(
encoder : @codec.Encoder,
metadata : Map[String, Bytes],
) -> Unit raise OcfError {
if !metadata.is_empty() {
encoder.write_long(Int64::from_int(metadata.length())) catch {
error => raise Datum(error)
}
for key, value in metadata {
encoder.write_string(key) catch {
error => raise Datum(error)
}
encoder.write_bytes(value) catch {
error => raise Datum(error)
}
}
}
encoder.write_long(0L) catch {
error => raise Datum(error)
}
}
///|
/// Serialize an OCF header (magic, metadata and 16-byte sync marker).
pub fn encode_header(header : OcfHeader) -> Bytes raise OcfError {
let encoder = @codec.Encoder::new()
encoder.write_raw(magic()) catch {
error => raise Datum(error)
}
write_metadata(encoder, header_metadata(header))
encoder.write_raw(header.sync_marker()) catch {
error => raise Datum(error)
}
encoder.to_bytes()
}
///|
/// Read one bounded Avro metadata map from a container stream.
fn read_metadata(
decoder : @codec.Decoder,
limits : OcfLimits,
) -> Map[String, Bytes] raise OcfError {
let result : Map[String, Bytes] = Map([])
let start = decoder.offset()
let mut count = decoder.read_long() catch { error => raise Datum(error) }
while count != 0L {
let entries = if count < 0L {
let positive = -count
if positive > Int64::from_int(limits.max_metadata_entries) {
raise TooManyRecords(count=positive, limit=limits.max_metadata_entries)
}
let block_bytes = decoder.read_long() catch {
error => raise Datum(error)
}
if block_bytes < 0L ||
block_bytes > Int64::from_int(limits.max_metadata_bytes) {
raise BlockTooLarge(
length=block_bytes.to_int(),
limit=limits.max_metadata_bytes,
)
}
positive
} else {
count
}
if entries > Int64::from_int(limits.max_metadata_entries - result.length()) {
raise TooManyRecords(count=entries, limit=limits.max_metadata_entries)
}
for _ in 0.. raise Datum(error) }
if key.is_empty() {
raise InvalidMetadataKey(key)
}
result[key] = decoder.read_bytes() catch { error => raise Datum(error) }
let used = decoder.offset() - start
if used > limits.max_metadata_bytes {
raise BlockTooLarge(length=used, limit=limits.max_metadata_bytes)
}
}
count = decoder.read_long() catch { error => raise Datum(error) }
}
result
}
///|
/// Decode the header from `decoder`, leaving it positioned at the first block.
fn decode_header_from(
decoder : @codec.Decoder,
limits : OcfLimits,
) -> OcfHeader raise OcfError {
let prefix = decoder.read_raw(4) catch { error => raise Datum(error) }
if prefix != magic() {
raise InvalidMagic
}
let metadata = read_metadata(decoder, limits)
let schema_bytes = match metadata.get("avro.schema") {
Some(value) => value
None => raise InvalidMetadataKey("missing required avro.schema metadata")
}
let schema_text = @utf8.decode(schema_bytes) catch {
_ => raise InvalidMetadataKey("avro.schema is not UTF-8")
}
let schema = @schema.parse(schema_text) catch {
error => raise HeaderSchema(error)
}
let codec = match metadata.get("avro.codec") {
None => Null
Some(value) => {
let text = @utf8.decode(value) catch {
_ => raise InvalidMetadataKey("avro.codec is not UTF-8")
}
CompressionCodec::parse(text)
}
}
// The public header retains application metadata only. Required protocol
// metadata is reconstructed on re-encoding, so a decoded header is safely
// reusable with `encode_header` / `encode_container`.
metadata.remove("avro.schema")
metadata.remove("avro.codec")
let sync_marker = decoder.read_raw(16) catch { error => raise Datum(error) }
OcfHeader::new(schema, codec~, metadata~, sync_marker)
}
///|
/// Decode only an OCF header. `bytes_consumed` is the first block offset.
pub fn decode_header(
input : Bytes,
limits? : OcfLimits = default_limits(),
) -> (OcfHeader, Int) raise OcfError {
let codec_limits = @codec.CodecLimits::new(
max_bytes=input.length(),
max_collection_items=limits.max_metadata_entries,
)
let decoder = @codec.Decoder::new(input, limits=codec_limits) catch {
error => raise Datum(error)
}
let header = decode_header_from(decoder, limits)
(header, decoder.offset())
}
///|
/// Encode one compressed OCF data block, including its record count, encoded
/// byte length and trailing sync marker.
fn encode_data_block(
header : OcfHeader,
record_count : Int,
uncompressed : Bytes,
limits : OcfLimits,
) -> Bytes raise OcfError {
if record_count <= 0 || record_count > limits.max_records_per_block {
raise TooManyRecords(
count=Int64::from_int(record_count),
limit=limits.max_records_per_block,
)
}
if uncompressed.length() > limits.max_block_bytes {
raise BlockTooLarge(
length=uncompressed.length(),
limit=limits.max_block_bytes,
)
}
let payload = compress_block(header.codec(), uncompressed)
if payload.length() > limits.max_block_bytes {
raise BlockTooLarge(length=payload.length(), limit=limits.max_block_bytes)
}
let encoder = @codec.Encoder::new(
limits=@codec.CodecLimits::new(max_bytes=payload.length() + 64),
)
encoder.write_long(Int64::from_int(record_count)) catch {
error => raise Datum(error)
}
encoder.write_long(Int64::from_int(payload.length())) catch {
error => raise Datum(error)
}
encoder.write_raw(payload) catch {
error => raise Datum(error)
}
encoder.write_raw(header.sync_marker()) catch {
error => raise Datum(error)
}
encoder.to_bytes()
}
///|
/// Encode records into one or more independently compressed OCF data blocks.
fn encode_container_with_config(
header : OcfHeader,
records : Array[@codec.Datum],
config : OcfBlockConfig,
limits : OcfLimits,
) -> Bytes raise OcfError {
let encoded_header = encode_header(header)
if records.is_empty() {
return encoded_header
}
let effective_records = if config.max_records < limits.max_records_per_block {
config.max_records
} else {
limits.max_records_per_block
}
let effective_bytes = if config.target_bytes < limits.max_block_bytes {
config.target_bytes
} else {
limits.max_block_bytes
}
let blocks : Array[Bytes] = []
let mut payload_encoder = @codec.Encoder::new(
limits=@codec.CodecLimits::new(max_bytes=limits.max_block_bytes),
)
let mut record_count = 0
for record in records {
let record_encoder = @codec.Encoder::new(
limits=@codec.CodecLimits::new(max_bytes=limits.max_block_bytes),
)
@codec.encode_into(header.schema(), record, record_encoder) catch {
error => raise Datum(error)
}
let record_bytes = record_encoder.to_bytes()
if record_bytes.length() > limits.max_block_bytes {
raise BlockTooLarge(
length=record_bytes.length(),
limit=limits.max_block_bytes,
)
}
let exceeds_target = record_count > 0 &&
record_bytes.length() > effective_bytes - payload_encoder.length()
if record_count >= effective_records || exceeds_target {
blocks.push(
encode_data_block(
header,
record_count,
payload_encoder.to_bytes(),
limits,
),
)
payload_encoder = @codec.Encoder::new(
limits=@codec.CodecLimits::new(max_bytes=limits.max_block_bytes),
)
record_count = 0
}
payload_encoder.write_raw(record_bytes) catch {
error => raise Datum(error)
}
record_count += 1
}
if record_count > 0 {
blocks.push(
encode_data_block(
header,
record_count,
payload_encoder.to_bytes(),
limits,
),
)
}
let mut total_bytes = encoded_header.length()
for block in blocks {
total_bytes += block.length()
}
let output = @codec.Encoder::new(
limits=@codec.CodecLimits::new(max_bytes=total_bytes),
)
output.write_raw(encoded_header) catch {
error => raise Datum(error)
}
for block in blocks {
output.write_raw(block) catch {
error => raise Datum(error)
}
}
output.to_bytes()
}
///|
/// Encode one complete Avro Object Container File as a single data block.
/// Use `encode_container_blocks` when output should be partitioned for bounded
/// decompression and incremental interoperability.
pub fn encode_container(
header : OcfHeader,
records : Array[@codec.Datum],
limits? : OcfLimits = default_limits(),
) -> Bytes raise OcfError {
if records.length() > limits.max_records_per_block {
raise TooManyRecords(
count=Int64::from_int(records.length()),
limit=limits.max_records_per_block,
)
}
let single = OcfBlockConfig::{
max_records: limits.max_records_per_block,
target_bytes: limits.max_block_bytes,
}
encode_container_with_config(header, records, single, limits)
}
///|
/// Encode one complete OCF file using independently compressed data blocks.
/// Blocks are closed when either configured record count or target
/// uncompressed byte size is reached. Record order is preserved.
pub fn encode_container_blocks(
header : OcfHeader,
records : Array[@codec.Datum],
config? : OcfBlockConfig = default_block_config(),
limits? : OcfLimits = default_limits(),
) -> Bytes raise OcfError {
encode_container_with_config(header, records, config, limits)
}
///|
/// Decode an entire OCF file, checking block sizes, record counts, sync markers
/// and all datum payloads against the schema carried by its header.
pub fn decode_container(
input : Bytes,
limits? : OcfLimits = default_limits(),
) -> (OcfHeader, Array[@codec.Datum]) raise OcfError {
let codec_limits = @codec.CodecLimits::new(
max_bytes=input.length(),
max_collection_items=limits.max_records_per_block,
)
let decoder = @codec.Decoder::new(input, limits=codec_limits) catch {
error => raise Datum(error)
}
let header = decode_header_from(decoder, limits)
let records : Array[@codec.Datum] = []
let mut block_index = 0
while decoder.remaining() > 0 {
let count = decoder.read_long() catch { error => raise Datum(error) }
if count <= 0L || count > Int64::from_int(limits.max_records_per_block) {
raise InvalidBlockCount(count~)
}
let length = decoder.read_long() catch { error => raise Datum(error) }
if length < 0L || length > Int64::from_int(limits.max_block_bytes) {
raise BlockTooLarge(length=length.to_int(), limit=limits.max_block_bytes)
}
let compressed = decoder.read_raw(length.to_int()) catch {
error => raise Datum(error)
}
let payload = decompress_block(header.codec(), compressed, limits)
let payload_limits = @codec.CodecLimits::new(
max_bytes=payload.length(),
max_collection_items=limits.max_records_per_block,
)
let payload_decoder = @codec.Decoder::new(payload, limits=payload_limits) catch {
error => raise Datum(error)
}
for _ in 0.. raise Datum(error)
},
)
}
payload_decoder.ensure_finished() catch {
error => raise Datum(error)
}
let marker = decoder.read_raw(16) catch { error => raise Datum(error) }
if marker != header.sync_marker() {
raise SyncMismatch(block_index~)
}
block_index += 1
}
(header, records)
}
///|
/// Ordered in-memory OCF writer. Records are appended in schema order and
/// partitioned into bounded data blocks when `finish` is called.
pub struct OcfWriter {
header : OcfHeader
limits : OcfLimits
block_config : OcfBlockConfig
records : Array[@codec.Datum]
mut finished : Bool
} derive(Debug)
///|
pub fn OcfWriter::new(
header : OcfHeader,
limits? : OcfLimits = default_limits(),
block_config? : OcfBlockConfig = default_block_config(),
) -> OcfWriter {
{ header, limits, block_config, records: [], finished: false }
}
///|
pub fn OcfWriter::append(
self : OcfWriter,
record : @codec.Datum,
) -> Unit raise OcfError {
if self.finished {
raise Compression(codec="writer", message="writer is already finished")
}
self.records.push(record)
}
///|
pub fn OcfWriter::finish(self : OcfWriter) -> Bytes raise OcfError {
if self.finished {
raise Compression(codec="writer", message="writer is already finished")
}
self.finished = true
encode_container_blocks(
self.header,
self.records,
config=self.block_config,
limits=self.limits,
)
}
///|
pub fn OcfWriter::record_count(self : OcfWriter) -> Int {
self.records.length()
}
///|
/// Ordered in-memory OCF reader. `next` returns records in block/file order.
pub struct OcfReader {
header : OcfHeader
records : Array[@codec.Datum]
mut index : Int
} derive(Debug)
///|
pub fn OcfReader::open(
input : Bytes,
limits? : OcfLimits = default_limits(),
) -> OcfReader raise OcfError {
let (header, records) = decode_container(input, limits~)
{ header, records, index: 0 }
}
///|
pub fn OcfReader::header(self : OcfReader) -> OcfHeader {
self.header
}
///|
pub fn OcfReader::next(self : OcfReader) -> @codec.Datum? {
if self.index >= self.records.length() {
None
} else {
let value = self.records[self.index]
self.index += 1
Some(value)
}
}
///|
pub fn OcfReader::remaining(self : OcfReader) -> Int {
self.records.length() - self.index
}