Skip to main content

kafrust_protocol/api/
describe_cluster.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5/// Kafka DescribeCluster API key.
6pub const API_KEY: i16 = 60;
7
8/// Describe broker endpoints.
9pub const BROKER_ENDPOINT_TYPE: i8 = 1;
10/// Describe controller endpoints.
11pub const CONTROLLER_ENDPOINT_TYPE: i8 = 2;
12
13/// Kafka DescribeCluster request for versions 0 and 1.
14#[derive(Debug, Clone, PartialEq, Eq)]
15pub struct DescribeClusterRequest {
16    pub correlation_id: i32,
17    pub client_id: Option<String>,
18    pub api_version: i16,
19    pub include_cluster_authorized_operations: bool,
20    pub endpoint_type: i8,
21}
22
23impl DescribeClusterRequest {
24    /// Encodes the flexible request header and body for the selected version.
25    pub fn encode(&self) -> Result<Vec<u8>> {
26        let mut encoder = Encoder::new();
27        RequestHeader {
28            api_key: API_KEY,
29            api_version: self.api_version,
30            correlation_id: self.correlation_id,
31            client_id: self.client_id.clone(),
32        }
33        .encode_v2(&mut encoder)?;
34        encoder.write_bool(self.include_cluster_authorized_operations);
35        if self.api_version >= 1 {
36            encoder.write_i8(self.endpoint_type);
37        }
38        encoder.write_empty_tagged_fields();
39        Ok(encoder.into_bytes())
40    }
41}
42
43/// One broker endpoint returned by DescribeCluster.
44#[derive(Debug, Clone, PartialEq, Eq)]
45pub struct DescribeClusterBroker {
46    pub node_id: i32,
47    pub host: String,
48    pub port: i32,
49    pub rack: Option<String>,
50}
51
52/// Kafka DescribeCluster response for versions 0 and 1.
53#[derive(Debug, Clone, PartialEq, Eq)]
54pub struct DescribeClusterResponse {
55    pub api_version: i16,
56    pub throttle_time_ms: i32,
57    pub error_code: i16,
58    pub error_message: Option<String>,
59    pub endpoint_type: Option<i8>,
60    pub cluster_id: String,
61    pub controller_id: i32,
62    pub brokers: Vec<DescribeClusterBroker>,
63    pub cluster_authorized_operations: i32,
64}
65
66impl DescribeClusterResponse {
67    /// Decodes a flexible response body for the selected version.
68    pub fn decode_body(decoder: &mut Decoder<'_>, api_version: i16) -> Result<Self> {
69        let throttle_time_ms = decoder.read_i32()?;
70        let error_code = decoder.read_i16()?;
71        let error_message = decoder.read_compact_nullable_string()?;
72        let endpoint_type = (api_version >= 1).then(|| decoder.read_i8()).transpose()?;
73        let cluster_id = decoder.read_compact_string()?;
74        let controller_id = decoder.read_i32()?;
75        let brokers = decoder
76            .read_compact_array("describe cluster brokers", |decoder| {
77                let node_id = decoder.read_i32()?;
78                let host = decoder.read_compact_string()?;
79                let port = decoder.read_i32()?;
80                let rack = decoder.read_compact_nullable_string()?;
81                decoder.read_tagged_fields()?;
82                Ok(DescribeClusterBroker {
83                    node_id,
84                    host,
85                    port,
86                    rack,
87                })
88            })?
89            .unwrap_or_default();
90        let cluster_authorized_operations = decoder.read_i32()?;
91        decoder.read_tagged_fields()?;
92        Ok(Self {
93            api_version,
94            throttle_time_ms,
95            error_code,
96            error_message,
97            endpoint_type,
98            cluster_id,
99            controller_id,
100            brokers,
101            cluster_authorized_operations,
102        })
103    }
104}
105
106#[cfg(test)]
107#[allow(clippy::unwrap_used)]
108mod tests {
109    use super::{DescribeClusterRequest, DescribeClusterResponse, API_KEY};
110    use crate::codec::{Decoder, Encoder};
111
112    #[test]
113    fn encodes_describe_cluster_v1_request() {
114        let request = DescribeClusterRequest {
115            correlation_id: 7,
116            client_id: Some("kafrust".to_owned()),
117            api_version: 1,
118            include_cluster_authorized_operations: true,
119            endpoint_type: 2,
120        };
121
122        assert_eq!(
123            request.encode().unwrap(),
124            [
125                0, 60, // API key
126                0, 1, // API version
127                0, 0, 0, 7, // correlation ID
128                0, 7, b'k', b'a', b'f', b'r', b'u', b's', b't',
129                0, // request header tagged fields
130                1, // include cluster authorized operations
131                2, // controller endpoint
132                0, // request tagged fields
133            ]
134        );
135        assert_eq!(API_KEY, 60);
136    }
137
138    #[test]
139    fn decodes_describe_cluster_v1_response() {
140        let mut body = Encoder::new();
141        body.write_i32(11);
142        body.write_i16(0);
143        body.write_compact_nullable_string(None).unwrap();
144        body.write_i8(2);
145        body.write_compact_string("cluster").unwrap();
146        body.write_i32(2);
147        body.write_compact_array(Some(&["broker-a".to_owned()]), |encoder, host| {
148            encoder.write_i32(1);
149            encoder.write_compact_string(host)?;
150            encoder.write_i32(9092);
151            encoder.write_compact_nullable_string(Some("rack-a"))?;
152            encoder.write_empty_tagged_fields();
153            Ok(())
154        })
155        .unwrap();
156        body.write_i32(0x1234);
157        body.write_empty_tagged_fields();
158
159        let bytes = body.into_bytes();
160        let mut decoder = Decoder::new(&bytes);
161        let response = DescribeClusterResponse::decode_body(&mut decoder, 1).unwrap();
162
163        assert_eq!(response.api_version, 1);
164        assert_eq!(response.throttle_time_ms, 11);
165        assert_eq!(response.endpoint_type, Some(2));
166        assert_eq!(response.cluster_id, "cluster");
167        assert_eq!(response.controller_id, 2);
168        assert_eq!(response.brokers.len(), 1);
169        assert_eq!(response.brokers[0].host, "broker-a");
170        assert_eq!(response.brokers[0].rack.as_deref(), Some("rack-a"));
171        assert_eq!(response.cluster_authorized_operations, 0x1234);
172        assert!(decoder.is_empty());
173    }
174}