Skip to main content

kafrust_protocol/api/
telemetry.rs

1//! Kafka client telemetry protocol types from KIP-714.
2
3use crate::codec::{Decoder, Encoder};
4use crate::error::Result;
5use crate::header::RequestHeader;
6
7/// Kafka API key for `GetTelemetrySubscriptions`.
8pub const GET_TELEMETRY_SUBSCRIPTIONS_API_KEY: i16 = 71;
9/// Kafka API key for `PushTelemetry`.
10pub const PUSH_TELEMETRY_API_KEY: i16 = 72;
11
12/// KIP-714 telemetry subscription request.
13#[derive(Debug, Clone, PartialEq, Eq)]
14pub struct GetTelemetrySubscriptionsRequestV0 {
15    pub correlation_id: i32,
16    pub client_id: Option<String>,
17    /// The all-zero UUID requests a new broker-assigned client instance ID.
18    pub client_instance_id: [u8; 16],
19}
20
21impl GetTelemetrySubscriptionsRequestV0 {
22    /// Encodes the flexible v0 request, including its request header.
23    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/// KIP-714 telemetry subscription response.
39#[derive(Debug, Clone, PartialEq, Eq)]
40pub struct GetTelemetrySubscriptionsResponseV0 {
41    pub throttle_time_ms: i32,
42    pub error_code: i16,
43    /// A non-zero value is returned when the request supplied the zero UUID.
44    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    /// Decodes the flexible response body after the response header.
55    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/// KIP-714 telemetry payload request.
89#[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    /// Compression type identifier from Kafka's record-batch codec values.
97    pub compression_type: i8,
98    /// OpenTelemetry MetricsData v1 protobuf bytes, optionally compressed as
99    /// declared by `compression_type`.
100    pub metrics: Vec<u8>,
101}
102
103impl PushTelemetryRequestV0 {
104    /// Encodes the flexible v0 request, including its request header.
105    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/// KIP-714 telemetry payload response.
125#[derive(Debug, Clone, PartialEq, Eq)]
126pub struct PushTelemetryResponseV0 {
127    pub throttle_time_ms: i32,
128    pub error_code: i16,
129}
130
131impl PushTelemetryResponseV0 {
132    /// Decodes the flexible response body after the response header.
133    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}