1use crate::codec::{Decoder, Encoder};
4use crate::error::Result;
5use crate::header::RequestHeader;
6
7pub const GET_TELEMETRY_SUBSCRIPTIONS_API_KEY: i16 = 71;
9pub const PUSH_TELEMETRY_API_KEY: i16 = 72;
11
12#[derive(Debug, Clone, PartialEq, Eq)]
14pub struct GetTelemetrySubscriptionsRequestV0 {
15 pub correlation_id: i32,
16 pub client_id: Option<String>,
17 pub client_instance_id: [u8; 16],
19}
20
21impl GetTelemetrySubscriptionsRequestV0 {
22 pub fn encode(&self) -> Result<Vec<u8>> {
24 let mut encoder = Encoder::new();
25 RequestHeader {
26 api_key: GET_TELEMETRY_SUBSCRIPTIONS_API_KEY,
27 api_version: 0,
28 correlation_id: self.correlation_id,
29 client_id: self.client_id.clone(),
30 }
31 .encode_v2(&mut encoder)?;
32 encoder.write_uuid(&self.client_instance_id);
33 encoder.write_empty_tagged_fields();
34 Ok(encoder.into_bytes())
35 }
36}
37
38#[derive(Debug, Clone, PartialEq, Eq)]
40pub struct GetTelemetrySubscriptionsResponseV0 {
41 pub throttle_time_ms: i32,
42 pub error_code: i16,
43 pub client_instance_id: [u8; 16],
45 pub subscription_id: i32,
46 pub accepted_compression_types: Vec<i8>,
47 pub push_interval_ms: i32,
48 pub telemetry_max_bytes: i32,
49 pub delta_temporality: bool,
50 pub requested_metrics: Vec<String>,
51}
52
53impl GetTelemetrySubscriptionsResponseV0 {
54 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
56 let throttle_time_ms = decoder.read_i32()?;
57 let error_code = decoder.read_i16()?;
58 let client_instance_id = decoder.read_uuid()?;
59 let subscription_id = decoder.read_i32()?;
60 let accepted_compression_types = decoder
61 .read_compact_array("accepted telemetry compression types", |decoder| {
62 decoder.read_i8()
63 })?
64 .unwrap_or_default();
65 let push_interval_ms = decoder.read_i32()?;
66 let telemetry_max_bytes = decoder.read_i32()?;
67 let delta_temporality = decoder.read_bool()?;
68 let requested_metrics = decoder
69 .read_compact_array("requested telemetry metrics", |decoder| {
70 decoder.read_compact_string()
71 })?
72 .unwrap_or_default();
73 decoder.read_tagged_fields()?;
74 Ok(Self {
75 throttle_time_ms,
76 error_code,
77 client_instance_id,
78 subscription_id,
79 accepted_compression_types,
80 push_interval_ms,
81 telemetry_max_bytes,
82 delta_temporality,
83 requested_metrics,
84 })
85 }
86}
87
88#[derive(Debug, Clone, PartialEq, Eq)]
90pub struct PushTelemetryRequestV0 {
91 pub correlation_id: i32,
92 pub client_id: Option<String>,
93 pub client_instance_id: [u8; 16],
94 pub subscription_id: i32,
95 pub terminating: bool,
96 pub compression_type: i8,
98 pub metrics: Vec<u8>,
101}
102
103impl PushTelemetryRequestV0 {
104 pub fn encode(&self) -> Result<Vec<u8>> {
106 let mut encoder = Encoder::new();
107 RequestHeader {
108 api_key: PUSH_TELEMETRY_API_KEY,
109 api_version: 0,
110 correlation_id: self.correlation_id,
111 client_id: self.client_id.clone(),
112 }
113 .encode_v2(&mut encoder)?;
114 encoder.write_uuid(&self.client_instance_id);
115 encoder.write_i32(self.subscription_id);
116 encoder.write_bool(self.terminating);
117 encoder.write_i8(self.compression_type);
118 encoder.write_compact_bytes(&self.metrics)?;
119 encoder.write_empty_tagged_fields();
120 Ok(encoder.into_bytes())
121 }
122}
123
124#[derive(Debug, Clone, PartialEq, Eq)]
126pub struct PushTelemetryResponseV0 {
127 pub throttle_time_ms: i32,
128 pub error_code: i16,
129}
130
131impl PushTelemetryResponseV0 {
132 pub fn decode_body(decoder: &mut Decoder<'_>) -> Result<Self> {
134 let throttle_time_ms = decoder.read_i32()?;
135 let error_code = decoder.read_i16()?;
136 decoder.read_tagged_fields()?;
137 Ok(Self {
138 throttle_time_ms,
139 error_code,
140 })
141 }
142}
143
144#[cfg(test)]
145#[allow(clippy::unwrap_used)]
146mod tests {
147 use super::*;
148
149 #[test]
150 fn encodes_get_telemetry_subscriptions_request() {
151 let request = GetTelemetrySubscriptionsRequestV0 {
152 correlation_id: 12,
153 client_id: Some("kafrust".to_owned()),
154 client_instance_id: [0; 16],
155 };
156 let encoded = request.encode().unwrap();
157 let mut decoder = Decoder::new(&encoded);
158 assert_eq!(
159 decoder.read_i16().unwrap(),
160 GET_TELEMETRY_SUBSCRIPTIONS_API_KEY
161 );
162 assert_eq!(decoder.read_i16().unwrap(), 0);
163 assert_eq!(decoder.read_i32().unwrap(), 12);
164 assert_eq!(
165 decoder.read_nullable_string().unwrap(),
166 Some("kafrust".to_owned())
167 );
168 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
169 assert_eq!(decoder.read_uuid().unwrap(), [0; 16]);
170 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
171 assert!(decoder.is_empty());
172 }
173
174 #[test]
175 fn decodes_telemetry_subscription_response() {
176 let mut encoder = Encoder::new();
177 encoder.write_i32(4);
178 encoder.write_i16(0);
179 encoder.write_uuid(&[1; 16]);
180 encoder.write_i32(7);
181 encoder
182 .write_compact_array(Some(&[0_i8, 2_i8]), |encoder, value| {
183 encoder.write_i8(*value);
184 Ok(())
185 })
186 .unwrap();
187 encoder.write_i32(30_000);
188 encoder.write_i32(1024 * 1024);
189 encoder.write_bool(true);
190 encoder
191 .write_compact_array(Some(&["org.apache.kafka.".to_owned()]), |encoder, value| {
192 encoder.write_compact_string(value)
193 })
194 .unwrap();
195 encoder.write_empty_tagged_fields();
196
197 let response = GetTelemetrySubscriptionsResponseV0::decode_body(&mut Decoder::new(
198 &encoder.into_bytes(),
199 ))
200 .unwrap();
201 assert_eq!(response.throttle_time_ms, 4);
202 assert_eq!(response.client_instance_id, [1; 16]);
203 assert_eq!(response.subscription_id, 7);
204 assert_eq!(response.accepted_compression_types, vec![0, 2]);
205 assert_eq!(response.push_interval_ms, 30_000);
206 assert_eq!(response.telemetry_max_bytes, 1024 * 1024);
207 assert!(response.delta_temporality);
208 assert_eq!(
209 response.requested_metrics,
210 vec!["org.apache.kafka.".to_owned()]
211 );
212 }
213
214 #[test]
215 fn encodes_push_telemetry_request_with_payload() {
216 let request = PushTelemetryRequestV0 {
217 correlation_id: 13,
218 client_id: None,
219 client_instance_id: [2; 16],
220 subscription_id: 7,
221 terminating: false,
222 compression_type: 0,
223 metrics: vec![1, 2, 3],
224 };
225 let encoded = request.encode().unwrap();
226 let mut decoder = Decoder::new(&encoded);
227 assert_eq!(decoder.read_i16().unwrap(), PUSH_TELEMETRY_API_KEY);
228 assert_eq!(decoder.read_i16().unwrap(), 0);
229 assert_eq!(decoder.read_i32().unwrap(), 13);
230 assert_eq!(decoder.read_nullable_string().unwrap(), None);
231 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
232 assert_eq!(decoder.read_uuid().unwrap(), [2; 16]);
233 assert_eq!(decoder.read_i32().unwrap(), 7);
234 assert!(!decoder.read_bool().unwrap());
235 assert_eq!(decoder.read_i8().unwrap(), 0);
236 assert_eq!(decoder.read_compact_bytes().unwrap(), vec![1, 2, 3]);
237 assert_eq!(decoder.read_tagged_fields().unwrap(), Vec::new());
238 assert!(decoder.is_empty());
239 }
240
241 #[test]
242 fn decodes_push_telemetry_response() {
243 let mut encoder = Encoder::new();
244 encoder.write_i32(9);
245 encoder.write_i16(0);
246 encoder.write_empty_tagged_fields();
247 let response =
248 PushTelemetryResponseV0::decode_body(&mut Decoder::new(&encoder.into_bytes())).unwrap();
249 assert_eq!(response.throttle_time_ms, 9);
250 assert_eq!(response.error_code, 0);
251 }
252}