Skip to main content

kafrust_protocol/api/
delegation_token.rs

1use crate::codec::{Decoder, Encoder};
2use crate::error::Result;
3use crate::header::RequestHeader;
4
5pub const CREATE_API_KEY: i16 = 38;
6pub const RENEW_API_KEY: i16 = 39;
7pub const EXPIRE_API_KEY: i16 = 40;
8pub const DESCRIBE_API_KEY: i16 = 41;
9
10#[derive(Debug, Clone, PartialEq, Eq)]
11pub struct DelegationTokenPrincipal {
12    pub principal_type: String,
13    pub principal_name: String,
14}
15
16#[derive(Debug, Clone, PartialEq, Eq)]
17pub struct CreateDelegationTokenRequest {
18    pub correlation_id: i32,
19    pub client_id: Option<String>,
20    pub owner: Option<DelegationTokenPrincipal>,
21    pub renewers: Vec<DelegationTokenPrincipal>,
22    pub max_lifetime_ms: i64,
23}
24
25impl CreateDelegationTokenRequest {
26    pub fn encode_v1(&self) -> Result<Vec<u8>> {
27        let mut encoder = Encoder::new();
28        RequestHeader {
29            api_key: CREATE_API_KEY,
30            api_version: 1,
31            correlation_id: self.correlation_id,
32            client_id: self.client_id.clone(),
33        }
34        .encode_v1(&mut encoder)?;
35        encoder.write_array(Some(&self.renewers), encode_legacy_principal)?;
36        encoder.write_i64(self.max_lifetime_ms);
37        Ok(encoder.into_bytes())
38    }
39
40    pub fn encode_v2(&self, api_version: i16) -> Result<Vec<u8>> {
41        let mut encoder = Encoder::new();
42        RequestHeader {
43            api_key: CREATE_API_KEY,
44            api_version,
45            correlation_id: self.correlation_id,
46            client_id: self.client_id.clone(),
47        }
48        .encode_v2(&mut encoder)?;
49        if api_version >= 3 {
50            let owner = self.owner.as_ref();
51            encoder
52                .write_compact_nullable_string(owner.map(|owner| owner.principal_type.as_str()))?;
53            encoder
54                .write_compact_nullable_string(owner.map(|owner| owner.principal_name.as_str()))?;
55        }
56        encoder.write_compact_array(Some(&self.renewers), encode_flexible_principal)?;
57        encoder.write_i64(self.max_lifetime_ms);
58        encoder.write_empty_tagged_fields();
59        Ok(encoder.into_bytes())
60    }
61}
62
63#[derive(Debug, Clone, PartialEq, Eq)]
64pub struct CreateDelegationTokenResponse {
65    pub error_code: i16,
66    pub owner: DelegationTokenPrincipal,
67    pub requester: Option<DelegationTokenPrincipal>,
68    pub issue_timestamp_ms: i64,
69    pub expiry_timestamp_ms: i64,
70    pub max_timestamp_ms: i64,
71    pub token_id: String,
72    pub hmac: Vec<u8>,
73    pub throttle_time_ms: i32,
74}
75
76impl CreateDelegationTokenResponse {
77    pub fn decode_body_v1(decoder: &mut Decoder<'_>) -> Result<Self> {
78        let error_code = decoder.read_i16()?;
79        let owner = decode_legacy_principal(decoder)?;
80        let issue_timestamp_ms = decoder.read_i64()?;
81        let expiry_timestamp_ms = decoder.read_i64()?;
82        let max_timestamp_ms = decoder.read_i64()?;
83        let token_id = decoder.read_string()?;
84        let hmac = decoder.read_bytes()?;
85        let throttle_time_ms = decoder.read_i32()?;
86        Ok(Self {
87            error_code,
88            owner,
89            requester: None,
90            issue_timestamp_ms,
91            expiry_timestamp_ms,
92            max_timestamp_ms,
93            token_id,
94            hmac,
95            throttle_time_ms,
96        })
97    }
98
99    pub fn decode_body_v2(decoder: &mut Decoder<'_>, api_version: i16) -> Result<Self> {
100        let error_code = decoder.read_i16()?;
101        let owner = decode_flexible_principal(decoder)?;
102        let requester = if api_version >= 3 {
103            Some(decode_flexible_principal(decoder)?)
104        } else {
105            None
106        };
107        let issue_timestamp_ms = decoder.read_i64()?;
108        let expiry_timestamp_ms = decoder.read_i64()?;
109        let max_timestamp_ms = decoder.read_i64()?;
110        let token_id = decoder.read_compact_string()?;
111        let hmac = decoder.read_compact_bytes()?;
112        let throttle_time_ms = decoder.read_i32()?;
113        decoder.read_tagged_fields()?;
114        Ok(Self {
115            error_code,
116            owner,
117            requester,
118            issue_timestamp_ms,
119            expiry_timestamp_ms,
120            max_timestamp_ms,
121            token_id,
122            hmac,
123            throttle_time_ms,
124        })
125    }
126}
127
128#[derive(Debug, Clone, PartialEq, Eq)]
129pub struct RenewDelegationTokenRequest {
130    pub correlation_id: i32,
131    pub client_id: Option<String>,
132    pub hmac: Vec<u8>,
133    pub renew_period_ms: i64,
134}
135
136impl RenewDelegationTokenRequest {
137    pub fn encode_v1(&self) -> Result<Vec<u8>> {
138        let mut encoder = Encoder::new();
139        RequestHeader {
140            api_key: RENEW_API_KEY,
141            api_version: 1,
142            correlation_id: self.correlation_id,
143            client_id: self.client_id.clone(),
144        }
145        .encode_v1(&mut encoder)?;
146        encoder.write_bytes(&self.hmac)?;
147        encoder.write_i64(self.renew_period_ms);
148        Ok(encoder.into_bytes())
149    }
150
151    pub fn encode_v2(&self) -> Result<Vec<u8>> {
152        let mut encoder = Encoder::new();
153        RequestHeader {
154            api_key: RENEW_API_KEY,
155            api_version: 2,
156            correlation_id: self.correlation_id,
157            client_id: self.client_id.clone(),
158        }
159        .encode_v2(&mut encoder)?;
160        encoder.write_compact_bytes(&self.hmac)?;
161        encoder.write_i64(self.renew_period_ms);
162        encoder.write_empty_tagged_fields();
163        Ok(encoder.into_bytes())
164    }
165}
166
167#[derive(Debug, Clone, PartialEq, Eq)]
168pub struct DelegationTokenOperationResponse {
169    pub error_code: i16,
170    pub expiry_timestamp_ms: i64,
171    pub throttle_time_ms: i32,
172}
173
174pub type RenewDelegationTokenResponse = DelegationTokenOperationResponse;
175pub type ExpireDelegationTokenResponse = DelegationTokenOperationResponse;
176
177impl RenewDelegationTokenResponse {
178    pub fn decode_body_v1(decoder: &mut Decoder<'_>) -> Result<Self> {
179        decode_operation_response(decoder, false)
180    }
181
182    pub fn decode_body_v2(decoder: &mut Decoder<'_>) -> Result<Self> {
183        decode_operation_response(decoder, true)
184    }
185}
186
187#[derive(Debug, Clone, PartialEq, Eq)]
188pub struct ExpireDelegationTokenRequest {
189    pub correlation_id: i32,
190    pub client_id: Option<String>,
191    pub hmac: Vec<u8>,
192    pub expiry_time_period_ms: i64,
193}
194
195impl ExpireDelegationTokenRequest {
196    pub fn encode_v1(&self) -> Result<Vec<u8>> {
197        let mut encoder = Encoder::new();
198        RequestHeader {
199            api_key: EXPIRE_API_KEY,
200            api_version: 1,
201            correlation_id: self.correlation_id,
202            client_id: self.client_id.clone(),
203        }
204        .encode_v1(&mut encoder)?;
205        encoder.write_bytes(&self.hmac)?;
206        encoder.write_i64(self.expiry_time_period_ms);
207        Ok(encoder.into_bytes())
208    }
209
210    pub fn encode_v2(&self) -> Result<Vec<u8>> {
211        let mut encoder = Encoder::new();
212        RequestHeader {
213            api_key: EXPIRE_API_KEY,
214            api_version: 2,
215            correlation_id: self.correlation_id,
216            client_id: self.client_id.clone(),
217        }
218        .encode_v2(&mut encoder)?;
219        encoder.write_compact_bytes(&self.hmac)?;
220        encoder.write_i64(self.expiry_time_period_ms);
221        encoder.write_empty_tagged_fields();
222        Ok(encoder.into_bytes())
223    }
224}
225
226#[derive(Debug, Clone, PartialEq, Eq)]
227pub struct DescribeDelegationTokenRequest {
228    pub correlation_id: i32,
229    pub client_id: Option<String>,
230    pub owners: Option<Vec<DelegationTokenPrincipal>>,
231}
232
233impl DescribeDelegationTokenRequest {
234    pub fn encode_v1(&self) -> Result<Vec<u8>> {
235        let mut encoder = Encoder::new();
236        RequestHeader {
237            api_key: DESCRIBE_API_KEY,
238            api_version: 1,
239            correlation_id: self.correlation_id,
240            client_id: self.client_id.clone(),
241        }
242        .encode_v1(&mut encoder)?;
243        encoder.write_array(self.owners.as_deref(), encode_legacy_principal)?;
244        Ok(encoder.into_bytes())
245    }
246
247    pub fn encode_v2(&self, api_version: i16) -> Result<Vec<u8>> {
248        let mut encoder = Encoder::new();
249        RequestHeader {
250            api_key: DESCRIBE_API_KEY,
251            api_version,
252            correlation_id: self.correlation_id,
253            client_id: self.client_id.clone(),
254        }
255        .encode_v2(&mut encoder)?;
256        encoder.write_compact_array(self.owners.as_deref(), encode_flexible_principal)?;
257        encoder.write_empty_tagged_fields();
258        Ok(encoder.into_bytes())
259    }
260}
261
262#[derive(Debug, Clone, PartialEq, Eq)]
263pub struct DescribeDelegationTokenResponse {
264    pub error_code: i16,
265    pub tokens: Vec<DescribedDelegationToken>,
266    pub throttle_time_ms: i32,
267}
268
269impl DescribeDelegationTokenResponse {
270    pub fn decode_body_v1(decoder: &mut Decoder<'_>) -> Result<Self> {
271        let error_code = decoder.read_i16()?;
272        let tokens = decoder
273            .read_array("described delegation tokens", decode_legacy_token)?
274            .unwrap_or_default();
275        let throttle_time_ms = decoder.read_i32()?;
276        Ok(Self {
277            error_code,
278            tokens,
279            throttle_time_ms,
280        })
281    }
282
283    pub fn decode_body_v2(decoder: &mut Decoder<'_>, api_version: i16) -> Result<Self> {
284        let error_code = decoder.read_i16()?;
285        let tokens = decoder
286            .read_compact_array("described delegation tokens", |decoder| {
287                decode_flexible_token(decoder, api_version)
288            })?
289            .unwrap_or_default();
290        let throttle_time_ms = decoder.read_i32()?;
291        decoder.read_tagged_fields()?;
292        Ok(Self {
293            error_code,
294            tokens,
295            throttle_time_ms,
296        })
297    }
298}
299
300#[derive(Debug, Clone, PartialEq, Eq)]
301pub struct DescribedDelegationToken {
302    pub owner: DelegationTokenPrincipal,
303    pub requester: Option<DelegationTokenPrincipal>,
304    pub issue_timestamp_ms: i64,
305    pub expiry_timestamp_ms: i64,
306    pub max_timestamp_ms: i64,
307    pub token_id: String,
308    pub hmac: Vec<u8>,
309    pub renewers: Vec<DelegationTokenPrincipal>,
310}
311
312fn decode_operation_response(
313    decoder: &mut Decoder<'_>,
314    flexible: bool,
315) -> Result<DelegationTokenOperationResponse> {
316    let response = DelegationTokenOperationResponse {
317        error_code: decoder.read_i16()?,
318        expiry_timestamp_ms: decoder.read_i64()?,
319        throttle_time_ms: decoder.read_i32()?,
320    };
321    if flexible {
322        decoder.read_tagged_fields()?;
323    }
324    Ok(response)
325}
326
327fn encode_legacy_principal(
328    encoder: &mut Encoder,
329    principal: &DelegationTokenPrincipal,
330) -> Result<()> {
331    encoder.write_string(&principal.principal_type)?;
332    encoder.write_string(&principal.principal_name)
333}
334
335fn encode_flexible_principal(
336    encoder: &mut Encoder,
337    principal: &DelegationTokenPrincipal,
338) -> Result<()> {
339    encoder.write_compact_string(&principal.principal_type)?;
340    encoder.write_compact_string(&principal.principal_name)?;
341    encoder.write_empty_tagged_fields();
342    Ok(())
343}
344
345fn decode_legacy_principal(decoder: &mut Decoder<'_>) -> Result<DelegationTokenPrincipal> {
346    Ok(DelegationTokenPrincipal {
347        principal_type: decoder.read_string()?,
348        principal_name: decoder.read_string()?,
349    })
350}
351
352fn decode_flexible_principal(decoder: &mut Decoder<'_>) -> Result<DelegationTokenPrincipal> {
353    Ok(DelegationTokenPrincipal {
354        principal_type: decoder.read_compact_string()?,
355        principal_name: decoder.read_compact_string()?,
356    })
357}
358
359fn decode_flexible_principal_struct(decoder: &mut Decoder<'_>) -> Result<DelegationTokenPrincipal> {
360    let principal = decode_flexible_principal(decoder)?;
361    decoder.read_tagged_fields()?;
362    Ok(principal)
363}
364
365fn decode_legacy_token(decoder: &mut Decoder<'_>) -> Result<DescribedDelegationToken> {
366    let owner = decode_legacy_principal(decoder)?;
367    let issue_timestamp_ms = decoder.read_i64()?;
368    let expiry_timestamp_ms = decoder.read_i64()?;
369    let max_timestamp_ms = decoder.read_i64()?;
370    let token_id = decoder.read_string()?;
371    let hmac = decoder.read_bytes()?;
372    let renewers = decoder
373        .read_array("delegation token renewers", decode_legacy_principal)?
374        .unwrap_or_default();
375    Ok(DescribedDelegationToken {
376        owner,
377        requester: None,
378        issue_timestamp_ms,
379        expiry_timestamp_ms,
380        max_timestamp_ms,
381        token_id,
382        hmac,
383        renewers,
384    })
385}
386
387fn decode_flexible_token(
388    decoder: &mut Decoder<'_>,
389    api_version: i16,
390) -> Result<DescribedDelegationToken> {
391    let owner = decode_flexible_principal(decoder)?;
392    let requester = if api_version >= 3 {
393        Some(decode_flexible_principal(decoder)?)
394    } else {
395        None
396    };
397    let issue_timestamp_ms = decoder.read_i64()?;
398    let expiry_timestamp_ms = decoder.read_i64()?;
399    let max_timestamp_ms = decoder.read_i64()?;
400    let token_id = decoder.read_compact_string()?;
401    let hmac = decoder.read_compact_bytes()?;
402    let renewers = decoder
403        .read_compact_array(
404            "delegation token renewers",
405            decode_flexible_principal_struct,
406        )?
407        .unwrap_or_default();
408    decoder.read_tagged_fields()?;
409    Ok(DescribedDelegationToken {
410        owner,
411        requester,
412        issue_timestamp_ms,
413        expiry_timestamp_ms,
414        max_timestamp_ms,
415        token_id,
416        hmac,
417        renewers,
418    })
419}
420
421#[cfg(test)]
422#[allow(clippy::unwrap_used)]
423mod tests {
424    use super::{
425        CreateDelegationTokenRequest, CreateDelegationTokenResponse, DelegationTokenPrincipal,
426        DescribeDelegationTokenRequest, DescribeDelegationTokenResponse,
427        ExpireDelegationTokenRequest, RenewDelegationTokenRequest, RenewDelegationTokenResponse,
428        CREATE_API_KEY, DESCRIBE_API_KEY,
429    };
430    use crate::codec::{Decoder, Encoder};
431
432    fn principal(kind: &str, name: &str) -> DelegationTokenPrincipal {
433        DelegationTokenPrincipal {
434            principal_type: kind.to_owned(),
435            principal_name: name.to_owned(),
436        }
437    }
438
439    #[test]
440    fn encodes_create_delegation_token_v1() {
441        let request = CreateDelegationTokenRequest {
442            correlation_id: 38,
443            client_id: Some("kafrust".to_owned()),
444            owner: None,
445            renewers: vec![principal("User", "alice")],
446            max_lifetime_ms: -1,
447        };
448        let bytes = request.encode_v1().unwrap();
449        assert_eq!(&bytes[0..4], &[0, CREATE_API_KEY as u8, 0, 1]);
450        assert!(bytes.windows(12).any(|window| {
451            window == [0, 4, b'U', b's', b'e', b'r', 0, 5, b'a', b'l', b'i', b'c']
452        }));
453        assert!(bytes.ends_with(&(-1_i64).to_be_bytes()));
454    }
455
456    #[test]
457    fn encodes_create_delegation_token_v3_with_owner() {
458        let request = CreateDelegationTokenRequest {
459            correlation_id: 39,
460            client_id: None,
461            owner: Some(principal("User", "owner")),
462            renewers: Vec::new(),
463            max_lifetime_ms: 60_000,
464        };
465        let bytes = request.encode_v2(3).unwrap();
466        assert_eq!(&bytes[0..4], &[0, CREATE_API_KEY as u8, 0, 3]);
467        assert!(bytes.ends_with(&[0]));
468    }
469
470    #[test]
471    fn decodes_create_delegation_token_v1_response() {
472        let mut bytes = Encoder::new();
473        bytes.write_i16(0);
474        bytes.write_string("User").unwrap();
475        bytes.write_string("alice").unwrap();
476        bytes.write_i64(10);
477        bytes.write_i64(20);
478        bytes.write_i64(30);
479        bytes.write_string("token-1").unwrap();
480        bytes.write_bytes(&[1, 2, 3]).unwrap();
481        bytes.write_i32(4);
482        let encoded = bytes.into_bytes();
483        let mut decoder = Decoder::new(&encoded);
484        let response = CreateDelegationTokenResponse::decode_body_v1(&mut decoder).unwrap();
485        assert_eq!(response.owner.principal_name, "alice");
486        assert_eq!(response.hmac, vec![1, 2, 3]);
487        assert_eq!(response.throttle_time_ms, 4);
488        assert!(decoder.is_empty());
489    }
490
491    #[test]
492    fn decodes_create_delegation_token_v3_response_without_nested_principal_tags() {
493        let mut bytes = Encoder::new();
494        bytes.write_i16(0);
495        bytes.write_compact_string("User").unwrap();
496        bytes.write_compact_string("owner").unwrap();
497        bytes.write_compact_string("User").unwrap();
498        bytes.write_compact_string("requester").unwrap();
499        bytes.write_i64(10);
500        bytes.write_i64(20);
501        bytes.write_i64(30);
502        bytes.write_compact_string("token-1").unwrap();
503        bytes.write_compact_bytes(b"secret-hmac").unwrap();
504        bytes.write_i32(4);
505        bytes.write_empty_tagged_fields();
506
507        let encoded = bytes.into_bytes();
508        let mut decoder = Decoder::new(&encoded);
509        let response = CreateDelegationTokenResponse::decode_body_v2(&mut decoder, 3).unwrap();
510        assert_eq!(response.owner.principal_name, "owner");
511        assert_eq!(
512            response.requester.as_ref().unwrap().principal_name,
513            "requester"
514        );
515        assert_eq!(response.hmac, b"secret-hmac");
516        assert_eq!(response.throttle_time_ms, 4);
517        assert!(decoder.is_empty());
518    }
519
520    #[test]
521    fn encodes_renew_and_expire_delegation_token_v2() {
522        let renew = RenewDelegationTokenRequest {
523            correlation_id: 40,
524            client_id: None,
525            hmac: vec![7, 8],
526            renew_period_ms: 100,
527        };
528        let expire = ExpireDelegationTokenRequest {
529            correlation_id: 41,
530            client_id: None,
531            hmac: vec![7, 8],
532            expiry_time_period_ms: 0,
533        };
534        assert_eq!(&renew.encode_v2().unwrap()[0..4], &[0, 39, 0, 2]);
535        assert_eq!(&expire.encode_v2().unwrap()[0..4], &[0, 40, 0, 2]);
536    }
537
538    #[test]
539    fn encodes_describe_delegation_tokens_v1_for_all_owners() {
540        let request = DescribeDelegationTokenRequest {
541            correlation_id: 42,
542            client_id: None,
543            owners: None,
544        };
545        let bytes = request.encode_v1().unwrap();
546        assert_eq!(&bytes[0..4], &[0, DESCRIBE_API_KEY as u8, 0, 1]);
547        assert_eq!(&bytes[bytes.len() - 4..], &(-1_i32).to_be_bytes());
548    }
549
550    #[test]
551    fn decodes_describe_delegation_tokens_v3_with_requester_and_tags() {
552        let mut bytes = Encoder::new();
553        bytes.write_i16(0);
554        bytes.write_unsigned_varint(2); // one token
555        bytes.write_compact_string("User").unwrap();
556        bytes.write_compact_string("alice").unwrap();
557        bytes.write_compact_string("User").unwrap();
558        bytes.write_compact_string("admin").unwrap();
559        bytes.write_i64(1);
560        bytes.write_i64(2);
561        bytes.write_i64(3);
562        bytes.write_compact_string("token-1").unwrap();
563        bytes.write_compact_bytes(&[9, 8]).unwrap();
564        bytes.write_unsigned_varint(2); // one renewer
565        bytes.write_compact_string("User").unwrap();
566        bytes.write_compact_string("renew").unwrap();
567        bytes.write_empty_tagged_fields();
568        bytes.write_empty_tagged_fields(); // token tags
569        bytes.write_i32(5);
570        bytes.write_empty_tagged_fields(); // response tags
571        let encoded = bytes.into_bytes();
572        let mut decoder = Decoder::new(&encoded);
573        let response = DescribeDelegationTokenResponse::decode_body_v2(&mut decoder, 3).unwrap();
574        assert_eq!(response.tokens[0].owner.principal_name, "alice");
575        assert_eq!(
576            response.tokens[0]
577                .requester
578                .as_ref()
579                .unwrap()
580                .principal_name,
581            "admin"
582        );
583        assert_eq!(response.tokens[0].hmac, vec![9, 8]);
584        assert_eq!(response.tokens[0].renewers[0].principal_name, "renew");
585        assert_eq!(response.throttle_time_ms, 5);
586        assert!(decoder.is_empty());
587    }
588
589    #[test]
590    fn decodes_renew_delegation_token_v2_response() {
591        let mut bytes = Encoder::new();
592        bytes.write_i16(0);
593        bytes.write_i64(99);
594        bytes.write_i32(3);
595        bytes.write_empty_tagged_fields();
596        let encoded = bytes.into_bytes();
597        let mut decoder = Decoder::new(&encoded);
598        let response = RenewDelegationTokenResponse::decode_body_v2(&mut decoder).unwrap();
599        assert_eq!(response.expiry_timestamp_ms, 99);
600        assert_eq!(response.throttle_time_ms, 3);
601        assert!(decoder.is_empty());
602    }
603}