use std::borrow::Cow;
use bytes::{Bytes, BytesMut};
use kacrab_protocol::{
KafkaString, KafkaUuid, Result,
generated::{
AddOffsetsToTxnRequestData, AddOffsetsToTxnResponseData, AddPartitionsToTxnRequestData,
AddPartitionsToTxnResponseData, AddRaftVoterRequestData, AddRaftVoterResponseData,
AlterClientQuotasRequestData, AlterClientQuotasResponseData, AlterConfigsRequestData,
AlterConfigsResponseData, AlterPartitionReassignmentsRequestData,
AlterPartitionReassignmentsResponseData, AlterReplicaLogDirsRequestData,
AlterReplicaLogDirsResponseData, AlterShareGroupOffsetsRequestData,
AlterShareGroupOffsetsResponseData, AlterUserScramCredentialsRequestData,
AlterUserScramCredentialsResponseData, ApiVersionsRequestData, ApiVersionsResponseData,
ConsumerGroupDescribeRequestData, ConsumerGroupDescribeResponseData,
ConsumerGroupHeartbeatRequestData, ConsumerGroupHeartbeatResponseData,
CreateAclsRequestData, CreateAclsResponseData, CreateDelegationTokenRequestData,
CreateDelegationTokenResponseData, CreatePartitionsRequestData,
CreatePartitionsResponseData, CreateTopicsRequestData, CreateTopicsResponseData,
DeleteAclsRequestData, DeleteAclsResponseData, DeleteGroupsRequestData,
DeleteGroupsResponseData, DeleteRecordsRequestData, DeleteRecordsResponseData,
DeleteShareGroupOffsetsRequestData, DeleteShareGroupOffsetsResponseData,
DeleteTopicsRequestData, DeleteTopicsResponseData, DescribeAclsRequestData,
DescribeAclsResponseData, DescribeClientQuotasRequestData,
DescribeClientQuotasResponseData, DescribeClusterRequestData, DescribeClusterResponseData,
DescribeConfigsRequestData, DescribeConfigsResponseData,
DescribeDelegationTokenRequestData, DescribeDelegationTokenResponseData,
DescribeGroupsRequestData, DescribeGroupsResponseData, DescribeLogDirsRequestData,
DescribeLogDirsResponseData, DescribeProducersRequestData, DescribeProducersResponseData,
DescribeQuorumRequestData, DescribeQuorumResponseData,
DescribeShareGroupOffsetsRequestData, DescribeShareGroupOffsetsResponseData,
DescribeTransactionsRequestData, DescribeTransactionsResponseData,
DescribeUserScramCredentialsRequestData, DescribeUserScramCredentialsResponseData,
ElectLeadersRequestData, ElectLeadersResponseData, EndTxnRequestData, EndTxnResponseData,
ExpireDelegationTokenRequestData, ExpireDelegationTokenResponseData, FetchRequestData,
FetchResponseData, FindCoordinatorRequestData, FindCoordinatorResponseData,
GetTelemetrySubscriptionsRequestData, GetTelemetrySubscriptionsResponseData,
HeartbeatRequestData, HeartbeatResponseData, IncrementalAlterConfigsRequestData,
IncrementalAlterConfigsResponseData, InitProducerIdRequestData, InitProducerIdResponseData,
JoinGroupRequestData, JoinGroupResponseData, LeaveGroupRequestData, LeaveGroupResponseData,
ListConfigResourcesRequestData, ListConfigResourcesResponseData, ListGroupsRequestData,
ListGroupsResponseData, ListOffsetsRequestData, ListOffsetsResponseData,
ListPartitionReassignmentsRequestData, ListPartitionReassignmentsResponseData,
ListTransactionsRequestData, ListTransactionsResponseData, MetadataRequestData,
MetadataResponseData, OffsetCommitRequestData, OffsetCommitResponseData,
OffsetDeleteRequestData, OffsetDeleteResponseData, OffsetFetchRequestData,
OffsetFetchResponseData, OffsetForLeaderEpochRequestData, OffsetForLeaderEpochResponseData,
ProduceRequestData, ProduceResponseData, PushTelemetryRequestData,
PushTelemetryResponseData, RemoveRaftVoterRequestData, RemoveRaftVoterResponseData,
RenewDelegationTokenRequestData, RenewDelegationTokenResponseData,
ShareGroupDescribeRequestData, ShareGroupDescribeResponseData,
StreamsGroupDescribeRequestData, StreamsGroupDescribeResponseData, SyncGroupRequestData,
SyncGroupResponseData, TxnOffsetCommitRequestData, TxnOffsetCommitResponseData,
UnregisterBrokerRequestData, UnregisterBrokerResponseData, UpdateFeaturesRequestData,
UpdateFeaturesResponseData, WriteTxnMarkersRequestData, WriteTxnMarkersResponseData,
},
};
pub trait RequestMessage {
fn write_request(&self, buf: &mut BytesMut, version: i16) -> Result<()>;
fn encoded_len(&self, version: i16) -> Result<usize>;
}
pub trait ResponseMessage: Sized {
fn read_response(buf: &mut Bytes, version: i16) -> Result<Self>;
}
macro_rules! impl_passthrough_message {
($($request:ty => $response:ty),+ $(,)?) => {
$(
impl RequestMessage for $request {
fn write_request(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
self.write(buf, version)?;
Ok(())
}
fn encoded_len(&self, version: i16) -> Result<usize> {
self.encoded_len(version)
}
}
impl ResponseMessage for $response {
fn read_response(buf: &mut Bytes, version: i16) -> Result<Self> {
Self::read(buf, version)
}
}
)+
};
}
impl_passthrough_message! {
ApiVersionsRequestData => ApiVersionsResponseData,
MetadataRequestData => MetadataResponseData,
InitProducerIdRequestData => InitProducerIdResponseData,
FindCoordinatorRequestData => FindCoordinatorResponseData,
AddPartitionsToTxnRequestData => AddPartitionsToTxnResponseData,
AddOffsetsToTxnRequestData => AddOffsetsToTxnResponseData,
TxnOffsetCommitRequestData => TxnOffsetCommitResponseData,
EndTxnRequestData => EndTxnResponseData,
GetTelemetrySubscriptionsRequestData => GetTelemetrySubscriptionsResponseData,
PushTelemetryRequestData => PushTelemetryResponseData,
}
impl RequestMessage for ProduceRequestData {
fn write_request(&self, buf: &mut BytesMut, version: i16) -> Result<()> {
normalize_produce_request(self, version).write(buf, version)?;
Ok(())
}
fn encoded_len(&self, version: i16) -> Result<usize> {
normalize_produce_request(self, version).encoded_len(version)
}
}
impl ResponseMessage for ProduceResponseData {
fn read_response(buf: &mut Bytes, version: i16) -> Result<Self> {
Self::read(buf, version)
}
}
impl_passthrough_message! {
CreateTopicsRequestData => CreateTopicsResponseData,
DeleteTopicsRequestData => DeleteTopicsResponseData,
CreatePartitionsRequestData => CreatePartitionsResponseData,
DescribeClusterRequestData => DescribeClusterResponseData,
DescribeConfigsRequestData => DescribeConfigsResponseData,
AlterConfigsRequestData => AlterConfigsResponseData,
ListGroupsRequestData => ListGroupsResponseData,
DescribeGroupsRequestData => DescribeGroupsResponseData,
DeleteGroupsRequestData => DeleteGroupsResponseData,
OffsetFetchRequestData => OffsetFetchResponseData,
OffsetCommitRequestData => OffsetCommitResponseData,
OffsetDeleteRequestData => OffsetDeleteResponseData,
IncrementalAlterConfigsRequestData => IncrementalAlterConfigsResponseData,
ElectLeadersRequestData => ElectLeadersResponseData,
ListOffsetsRequestData => ListOffsetsResponseData,
DeleteRecordsRequestData => DeleteRecordsResponseData,
DescribeProducersRequestData => DescribeProducersResponseData,
DescribeTransactionsRequestData => DescribeTransactionsResponseData,
ListTransactionsRequestData => ListTransactionsResponseData,
DescribeLogDirsRequestData => DescribeLogDirsResponseData,
AlterPartitionReassignmentsRequestData => AlterPartitionReassignmentsResponseData,
ListPartitionReassignmentsRequestData => ListPartitionReassignmentsResponseData,
UpdateFeaturesRequestData => UpdateFeaturesResponseData,
UnregisterBrokerRequestData => UnregisterBrokerResponseData,
DescribeAclsRequestData => DescribeAclsResponseData,
CreateAclsRequestData => CreateAclsResponseData,
DeleteAclsRequestData => DeleteAclsResponseData,
DescribeClientQuotasRequestData => DescribeClientQuotasResponseData,
AlterClientQuotasRequestData => AlterClientQuotasResponseData,
DescribeUserScramCredentialsRequestData => DescribeUserScramCredentialsResponseData,
AlterUserScramCredentialsRequestData => AlterUserScramCredentialsResponseData,
CreateDelegationTokenRequestData => CreateDelegationTokenResponseData,
RenewDelegationTokenRequestData => RenewDelegationTokenResponseData,
ExpireDelegationTokenRequestData => ExpireDelegationTokenResponseData,
DescribeDelegationTokenRequestData => DescribeDelegationTokenResponseData,
AlterReplicaLogDirsRequestData => AlterReplicaLogDirsResponseData,
WriteTxnMarkersRequestData => WriteTxnMarkersResponseData,
LeaveGroupRequestData => LeaveGroupResponseData,
ConsumerGroupDescribeRequestData => ConsumerGroupDescribeResponseData,
ListConfigResourcesRequestData => ListConfigResourcesResponseData,
DescribeQuorumRequestData => DescribeQuorumResponseData,
AddRaftVoterRequestData => AddRaftVoterResponseData,
RemoveRaftVoterRequestData => RemoveRaftVoterResponseData,
ShareGroupDescribeRequestData => ShareGroupDescribeResponseData,
StreamsGroupDescribeRequestData => StreamsGroupDescribeResponseData,
DescribeShareGroupOffsetsRequestData => DescribeShareGroupOffsetsResponseData,
AlterShareGroupOffsetsRequestData => AlterShareGroupOffsetsResponseData,
DeleteShareGroupOffsetsRequestData => DeleteShareGroupOffsetsResponseData,
}
impl_passthrough_message! {
FetchRequestData => FetchResponseData,
JoinGroupRequestData => JoinGroupResponseData,
SyncGroupRequestData => SyncGroupResponseData,
HeartbeatRequestData => HeartbeatResponseData,
OffsetForLeaderEpochRequestData => OffsetForLeaderEpochResponseData,
ConsumerGroupHeartbeatRequestData => ConsumerGroupHeartbeatResponseData,
}
fn normalize_produce_request(
request: &ProduceRequestData,
version: i16,
) -> Cow<'_, ProduceRequestData> {
let needs_clear = request.topic_data.iter().any(|topic| {
if version >= 13 {
topic.name != KafkaString::default()
} else {
topic.topic_id != KafkaUuid::ZERO
}
});
if !needs_clear {
return Cow::Borrowed(request);
}
let mut normalized = request.clone();
for topic in &mut normalized.topic_data {
if version >= 13 {
topic.name = KafkaString::default();
} else {
topic.topic_id = KafkaUuid::ZERO;
}
}
Cow::Owned(normalized)
}