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); 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); 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(); bytes.write_i32(5);
570 bytes.write_empty_tagged_fields(); 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}