datum-mq 0.10.5

Kafka sources and sinks for Datum streams, with native and rdkafka backends
Documentation
mod codec;
pub(crate) mod compression;
mod messages;
mod records;

pub(crate) use codec::{Decoder, Encoder};
pub(crate) use compression::CompressionCodec;
pub(crate) use messages::{
    ApiVersionsResponse, ConsumerGroupProtocol, FetchPartitionRequest, FetchResponse,
    FetchTopicRequest, FindCoordinatorResponse, HeartbeatResponse, JoinGroupMember,
    JoinGroupResponse, LeaveGroupResponse, ListOffsetPartitionRequest, ListOffsetResult,
    ListOffsetTopicRequest, MetadataResponse, OffsetCommitPartitionRequest,
    OffsetCommitTopicRequest, OffsetFetchResponse, OffsetFetchResult, OffsetFetchTopicRequest,
    ProducePartitionRequest, ProduceTopicRequest, SaslAuthenticateResponse, SaslHandshakeResponse,
    SyncGroupAssignment, SyncGroupResponse, decode_consumer_protocol_assignment,
    decode_consumer_protocol_subscription, decode_list_offsets_response,
    decode_offset_commit_response, decode_produce_response, encode_api_versions_request,
    encode_consumer_protocol_assignment, encode_consumer_protocol_subscription,
    encode_fetch_request, encode_find_coordinator_request, encode_heartbeat_request,
    encode_join_group_request, encode_leave_group_request, encode_list_offsets_request,
    encode_metadata_request, encode_offset_commit_request, encode_offset_fetch_request,
    encode_produce_request, encode_sasl_authenticate_request, encode_sasl_handshake_request,
    encode_sync_group_request,
};
#[cfg(test)]
pub(crate) use messages::{MetadataPartition, MetadataTopic};
pub(crate) use records::{decode_record_batches, encode_produce_record_batch};

pub(crate) const API_KEY_PRODUCE: i16 = 0;
pub(crate) const API_KEY_FETCH: i16 = 1;
pub(crate) const API_KEY_LIST_OFFSETS: i16 = 2;
pub(crate) const API_KEY_METADATA: i16 = 3;
pub(crate) const API_KEY_OFFSET_COMMIT: i16 = 8;
pub(crate) const API_KEY_OFFSET_FETCH: i16 = 9;
pub(crate) const API_KEY_FIND_COORDINATOR: i16 = 10;
pub(crate) const API_KEY_JOIN_GROUP: i16 = 11;
pub(crate) const API_KEY_HEARTBEAT: i16 = 12;
pub(crate) const API_KEY_LEAVE_GROUP: i16 = 13;
pub(crate) const API_KEY_SYNC_GROUP: i16 = 14;
pub(crate) const API_KEY_SASL_HANDSHAKE: i16 = 17;
pub(crate) const API_KEY_API_VERSIONS: i16 = 18;
pub(crate) const API_KEY_SASL_AUTHENTICATE: i16 = 36;

pub(crate) const API_VERSION_PRODUCE: i16 = 9;
pub(crate) const API_VERSION_FETCH: i16 = 12;
pub(crate) const API_VERSION_LIST_OFFSETS: i16 = 6;
pub(crate) const API_VERSION_METADATA: i16 = 9;
pub(crate) const API_VERSION_OFFSET_COMMIT: i16 = 8;
pub(crate) const API_VERSION_OFFSET_FETCH: i16 = 7;
pub(crate) const API_VERSION_FIND_COORDINATOR: i16 = 3;
pub(crate) const API_VERSION_JOIN_GROUP: i16 = 6;
pub(crate) const API_VERSION_HEARTBEAT: i16 = 4;
pub(crate) const API_VERSION_LEAVE_GROUP: i16 = 4;
pub(crate) const API_VERSION_SYNC_GROUP: i16 = 4;
pub(crate) const API_VERSION_API_VERSIONS: i16 = 3;
pub(crate) const API_VERSION_SASL_HANDSHAKE: i16 = 1;
pub(crate) const API_VERSION_SASL_AUTHENTICATE: i16 = 1;

pub(crate) fn request_header_version(api_version: i16) -> i16 {
    if api_version >= 3 { 2 } else { 1 }
}

pub(crate) fn response_header_version(api_key: i16, api_version: i16) -> i16 {
    if api_key == API_KEY_API_VERSIONS {
        0
    } else if is_flexible_response(api_key, api_version) {
        1
    } else {
        0
    }
}

fn is_flexible_response(api_key: i16, api_version: i16) -> bool {
    match api_key {
        API_KEY_FIND_COORDINATOR => api_version >= 3,
        API_KEY_JOIN_GROUP => api_version >= 6,
        API_KEY_HEARTBEAT => api_version >= 4,
        API_KEY_LEAVE_GROUP => api_version >= 4,
        API_KEY_SYNC_GROUP => api_version >= 4,
        API_KEY_OFFSET_COMMIT => api_version >= 8,
        API_KEY_OFFSET_FETCH => api_version >= 6,
        API_KEY_SASL_AUTHENTICATE => api_version >= 2,
        _ => api_version >= 6,
    }
}