///|
/// 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
}