kafkaesque 0.0.15

Kafka implementations in Rust (Unofficial)
Documentation
use crate::formats::api_keys::ApiKey;
use crate::formats::codec::{Read, Write};
use crate::formats::crc32::Crc32;
use crate::formats::request::{ApiVersion, RequestMessage};
use crate::formats::{ErrorCode, VarInt, VarIntArray, VarIntString};
use crate::formats::{FixedLength, Result};
use bytes::Bytes;
use futures::executor::block_on;
use tokio::io::{AsyncRead, AsyncWrite};

#[derive(Debug, Write, Read, RequestMessage)]
#[request_message(version = 0, key = "Produce")]
pub struct ProduceReqV0 {
    pub acks: i16,
    pub timeout_ms: i32,
}

#[derive(Debug, Write, Read)]
pub struct ProduceReqV0TopicData {
    pub topic_name: String,
    pub topic_data: Vec<ProduceReqV0PartitionData>,
}

#[derive(Debug, Write, Read)]
pub struct ProduceReqV0PartitionData {
    pub partition_id: i32,
    pub records: RecordBatch,
}

#[derive(Debug, Write, Read)]
pub struct ProduceRespV0 {}

#[derive(Debug, Clone)]
pub struct RecordBatchInput {
    pub base_offset: i64,
    pub batch_length: i32,
    pub partition_leader_epoch: i32,
    pub magic: i8,
    pub attributes: RecordBatchAttributes,
    pub last_offset_delta: i32,
    pub base_timestamp: i64,
    pub max_timestamp: i64,
    pub producer_id: i64,
    pub producer_epoch: i16,
    pub base_sequence: i32,
    pub records: Vec<Record>,
}

impl RecordBatchInput {
    fn into_record_batch(self) -> RecordBatch {
        let crc = {
            let mut crc = Crc32::default();
            let crc_input = RecordBatchCrcInput {
                attributes: self.attributes,
                last_offset_delta: self.last_offset_delta,
                base_timestamp: self.base_timestamp,
                max_timestamp: self.max_timestamp,
                producer_id: self.producer_id,
                producer_epoch: self.producer_epoch,
                base_sequence: self.base_sequence,
                records: &self.records,
            };
            block_on(crc_input.write_to(&mut crc)).expect("Failed to write to Crc32");
            crc.finalize()
        };
        RecordBatch {
            base_offset: self.base_offset,
            batch_length: self.batch_length,
            partition_leader_epoch: self.partition_leader_epoch,
            magic: self.magic,
            crc,
            attributes: self.attributes,
            last_offset_delta: self.last_offset_delta,
            base_timestamp: self.base_timestamp,
            max_timestamp: self.max_timestamp,
            producer_id: self.producer_id,
            producer_epoch: self.producer_epoch,
            base_sequence: self.base_sequence,
            records: self.records,
        }
    }
}

#[derive(Write, Debug)]
pub struct RecordBatchCrcInput<'a> {
    pub attributes: RecordBatchAttributes,
    pub last_offset_delta: i32,
    pub base_timestamp: i64,
    pub max_timestamp: i64,
    pub producer_id: i64,
    pub producer_epoch: i16,
    pub base_sequence: i32,
    pub records: &'a [Record],
}

#[derive(Debug, Clone, Write, Read)]
pub struct RecordBatch {
    pub base_offset: i64,
    pub batch_length: i32,
    pub partition_leader_epoch: i32,
    pub magic: i8,
    pub crc: u32,
    pub attributes: RecordBatchAttributes,
    pub last_offset_delta: i32,
    pub base_timestamp: i64,
    pub max_timestamp: i64,
    pub producer_id: i64,
    pub producer_epoch: i16,
    pub base_sequence: i32,
    pub records: Vec<Record>,
}

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RecordBatchCompression {
    None,
    Gzip,
    Snappy,
    Lz4,
    Zstd,
}

#[derive(Debug, Clone, Copy)]
pub struct RecordBatchAttributes {
    pub compression: RecordBatchCompression,
    pub is_transactional: bool,
    pub is_control_batch: bool,
    pub has_delete_horizons: bool,
}

impl RecordBatchAttributes {
    const COMPRESSION_GZIP: i16 = 0b0000000000000001;
    const COMPRESSION_SNAPPY: i16 = 0b0000000000000010;
    const COMPRESSION_LZ4: i16 = 0b0000000000000011;
    const COMPRESSION_ZSTD: i16 = 0b0000000000000110;
    const IS_TRANSACTIONAL: i16 = 0b0000000000001000;
    const IS_CONTROL_BATCH: i16 = 0b0000000000010000;
    const HAS_DELETE_HORIZONS: i16 = 0b0000000000100000;
}

impl Write for RecordBatchAttributes {
    fn calculate_size(&self) -> i32 {
        i16::SIZE
    }
    async fn write_to(&self, writer: &mut (dyn AsyncWrite + Send + Unpin)) -> Result<()> {
        let compression = match self.compression {
            RecordBatchCompression::None => 0,
            RecordBatchCompression::Gzip => Self::COMPRESSION_GZIP,
            RecordBatchCompression::Snappy => Self::COMPRESSION_SNAPPY,
            RecordBatchCompression::Lz4 => Self::COMPRESSION_LZ4,
            RecordBatchCompression::Zstd => Self::COMPRESSION_ZSTD,
        };
        let is_transactional = if self.is_transactional {
            Self::IS_TRANSACTIONAL
        } else {
            0
        };
        let is_control = if self.is_control_batch {
            Self::IS_CONTROL_BATCH
        } else {
            0
        };
        let has_delete_horizons = if self.has_delete_horizons {
            Self::HAS_DELETE_HORIZONS
        } else {
            0
        };
        let n: i16 = compression & is_transactional & is_control & has_delete_horizons;
        n.write_to(writer).await
    }
}

impl Read for RecordBatchAttributes {
    async fn read_from(reader: &mut (dyn AsyncRead + Send + Unpin)) -> Result<Self> {
        let n = i16::read_from(reader).await?;
        let compression = if (n & Self::COMPRESSION_GZIP != 0) {
            RecordBatchCompression::Gzip
        } else if (n & Self::COMPRESSION_SNAPPY != 0) {
            RecordBatchCompression::Snappy
        } else if (n & Self::COMPRESSION_LZ4 != 0) {
            RecordBatchCompression::Lz4
        } else if (n & Self::COMPRESSION_ZSTD != 0) {
            RecordBatchCompression::Zstd
        } else {
            RecordBatchCompression::None
        };
        let is_transactional = n & Self::IS_TRANSACTIONAL != 0;
        let is_control_batch = n & Self::IS_CONTROL_BATCH != 0;
        let has_delete_horizons = n & Self::HAS_DELETE_HORIZONS != 0;
        Ok(RecordBatchAttributes {
            compression,
            is_transactional,
            is_control_batch,
            has_delete_horizons,
        })
    }
}

#[derive(Debug, Clone, Write, Read)]
pub struct Record {
    pub length: VarInt,
    pub attributes: i8,
    pub timestamp_delta: VarInt,
    pub offset_delta: VarInt,
    pub key: VarIntArray<u8>,
    pub value: VarIntArray<u8>,
    pub headers: Vec<RecordHeader>,
}

#[derive(Debug, Clone, Write, Read)]
pub struct RecordHeader {
    pub key: VarIntString,
    pub value: VarIntArray<u8>,
}

#[derive(Debug, Clone, Write, Read)]
pub struct ControlRecord {
    pub version: i16,
    pub ty: i16,
}