1use std::collections::{BTreeMap, HashMap, HashSet};
16
17use common_time::util::current_time_millis;
18use derive_builder::Builder;
19use serde::ser::SerializeSeq;
20use serde::{Deserialize, Deserializer, Serialize, Serializer};
21use store_api::region_engine::RegionRole;
22use store_api::storage::{RegionId, RegionNumber};
23use strum::AsRefStr;
24
25use crate::DatanodeId;
26use crate::key::RegionDistribution;
27use crate::peer::Peer;
28
29pub fn region_distribution(region_routes: &[RegionRoute]) -> RegionDistribution {
34 let mut regions_id_map = RegionDistribution::new();
35 for route in region_routes.iter() {
36 if let Some(peer) = route.leader_peer.as_ref() {
37 let region_number = route.region.id.region_number();
38 regions_id_map
39 .entry(peer.id)
40 .or_default()
41 .add_leader_region(region_number);
42 }
43 for peer in route.follower_peers.iter() {
44 let region_number = route.region.id.region_number();
45 regions_id_map
46 .entry(peer.id)
47 .or_default()
48 .add_follower_region(region_number);
49 }
50 }
51 for (_, region_role_set) in regions_id_map.iter_mut() {
52 region_role_set.sort()
54 }
55 regions_id_map
56}
57
58pub fn find_leaders(region_routes: &[RegionRoute]) -> HashSet<Peer> {
60 region_routes
61 .iter()
62 .flat_map(|x| &x.leader_peer)
63 .cloned()
64 .collect()
65}
66
67pub fn find_followers(region_routes: &[RegionRoute]) -> HashSet<Peer> {
69 region_routes
70 .iter()
71 .flat_map(|x| &x.follower_peers)
72 .cloned()
73 .collect()
74}
75
76pub fn operating_leader_regions(region_routes: &[RegionRoute]) -> Vec<(RegionId, DatanodeId)> {
78 region_routes
79 .iter()
80 .filter_map(|route| {
81 route
82 .leader_peer
83 .as_ref()
84 .map(|leader| (route.region.id, leader.id))
85 })
86 .collect::<Vec<_>>()
87}
88
89pub fn operating_leader_region_roles(
91 region_routes: &[RegionRoute],
92) -> Vec<(RegionId, DatanodeId, RegionRole)> {
93 region_routes
94 .iter()
95 .filter_map(|route| {
96 let role = route.leader_region_role()?;
97 let leader = route.leader_peer.as_ref()?;
98 Some((route.region.id, leader.id, role))
99 })
100 .collect()
101}
102
103pub fn convert_to_region_leader_map(region_routes: &[RegionRoute]) -> HashMap<RegionNumber, &Peer> {
107 region_routes
108 .iter()
109 .filter_map(|x| {
110 x.leader_peer
111 .as_ref()
112 .map(|leader| (x.region.id.region_number(), leader))
113 })
114 .collect::<HashMap<_, _>>()
115}
116
117pub fn find_region_leader(
118 region_routes: &[RegionRoute],
119 region_number: RegionNumber,
120) -> Option<Peer> {
121 region_routes
122 .iter()
123 .find(|x| x.region.id.region_number() == region_number)
124 .and_then(|r| r.leader_peer.as_ref())
125 .cloned()
126}
127
128pub fn find_leader_regions(region_routes: &[RegionRoute], datanode: &Peer) -> Vec<RegionNumber> {
130 region_routes
131 .iter()
132 .filter_map(|x| {
133 if let Some(peer) = &x.leader_peer
134 && peer == datanode
135 {
136 return Some(x.region.id.region_number());
137 }
138 None
139 })
140 .collect()
141}
142
143pub fn find_follower_regions(region_routes: &[RegionRoute], datanode: &Peer) -> Vec<RegionNumber> {
145 region_routes
146 .iter()
147 .filter_map(|x| {
148 if x.follower_peers.contains(datanode) {
149 return Some(x.region.id.region_number());
150 }
151 None
152 })
153 .collect()
154}
155
156#[derive(Debug, Clone, Default, Deserialize, Serialize, PartialEq, Builder)]
157pub struct RegionRoute {
158 pub region: Region,
159 #[builder(setter(into, strip_option))]
160 pub leader_peer: Option<Peer>,
161 #[builder(setter(into), default)]
162 pub follower_peers: Vec<Peer>,
163 #[builder(setter(into, strip_option), default)]
165 #[serde(
166 default,
167 alias = "leader_status",
168 skip_serializing_if = "Option::is_none"
169 )]
170 pub leader_state: Option<LeaderState>,
171 #[serde(default)]
173 #[builder(default = "self.default_leader_down_since()")]
174 pub leader_down_since: Option<i64>,
175 #[builder(setter(into, strip_option), default)]
177 #[serde(default, skip_serializing_if = "Option::is_none")]
178 pub write_route_policy: Option<WriteRoutePolicy>,
179}
180
181impl RegionRouteBuilder {
182 fn default_leader_down_since(&self) -> Option<i64> {
183 match self.leader_state {
184 Some(Some(LeaderState::Downgrading)) => Some(current_time_millis()),
185 _ => None,
186 }
187 }
188}
189
190#[derive(Debug, Clone, Copy, Deserialize, Serialize, PartialEq, AsRefStr)]
193#[strum(serialize_all = "UPPERCASE")]
194pub enum LeaderState {
195 #[serde(alias = "Downgraded")]
200 Downgrading,
201 Staging,
206}
207
208#[derive(Debug, Clone, Copy, Deserialize, Serialize, PartialEq, Eq)]
210pub enum WriteRoutePolicy {
211 Normal,
213 IgnoreAllWrites,
219}
220
221impl RegionRoute {
222 pub fn is_ignore_all_writes(&self) -> bool {
224 matches!(
225 self.write_route_policy,
226 Some(WriteRoutePolicy::IgnoreAllWrites)
227 )
228 }
229
230 pub fn set_ignore_all_writes(&mut self) {
232 self.write_route_policy = Some(WriteRoutePolicy::IgnoreAllWrites);
233 }
234
235 pub fn clear_ignore_all_writes(&mut self) {
237 if self.write_route_policy == Some(WriteRoutePolicy::IgnoreAllWrites) {
238 self.write_route_policy = None;
239 }
240 }
241
242 pub fn is_leader_downgrading(&self) -> bool {
250 matches!(self.leader_state, Some(LeaderState::Downgrading))
251 }
252
253 pub fn is_leader_staging(&self) -> bool {
255 matches!(self.leader_state, Some(LeaderState::Staging))
256 }
257
258 pub fn leader_region_role(&self) -> Option<RegionRole> {
260 self.leader_peer.as_ref().map(|_| {
261 if self.is_leader_staging() {
262 RegionRole::StagingLeader
263 } else if self.is_leader_downgrading() {
264 RegionRole::DowngradingLeader
265 } else {
266 RegionRole::Leader
267 }
268 })
269 }
270
271 pub fn downgrade_leader(&mut self) {
283 self.leader_down_since = Some(current_time_millis());
284 self.leader_state = Some(LeaderState::Downgrading)
285 }
286
287 pub fn set_leader_staging(&mut self) {
289 self.leader_state = Some(LeaderState::Staging);
290 self.leader_down_since = None;
292 }
293
294 pub fn clear_leader_staging(&mut self) {
296 if self.leader_state == Some(LeaderState::Staging) {
297 self.leader_state = None;
298 self.leader_down_since = None;
299 }
300 }
301
302 pub fn leader_down_millis(&self) -> Option<i64> {
304 self.leader_down_since
305 .map(|start| current_time_millis() - start)
306 }
307
308 pub fn set_leader_state(&mut self, state: Option<LeaderState>) -> bool {
312 let updated = self.leader_state != state;
313
314 match (state, updated) {
315 (Some(LeaderState::Downgrading), true) => {
316 self.leader_down_since = Some(current_time_millis());
317 }
318 (Some(LeaderState::Downgrading), false) => {
319 }
321 _ => {
322 self.leader_down_since = None;
323 }
324 }
325
326 self.leader_state = state;
327 updated
328 }
329}
330
331#[derive(Debug, Clone, Default, PartialEq, Serialize)]
332pub struct Region {
333 pub id: RegionId,
334 pub name: String,
335 pub attrs: BTreeMap<String, String>,
336 pub partition_expr: String,
338}
339
340#[derive(Debug, Deserialize)]
341struct RegionDe {
342 id: RegionId,
343 name: String,
344 #[serde(default)]
345 attrs: BTreeMap<String, String>,
346 #[serde(default)]
347 partition: Option<LegacyPartition>,
348 #[serde(default)]
349 partition_expr: String,
350}
351
352impl<'de> Deserialize<'de> for Region {
353 fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
354 where
355 D: Deserializer<'de>,
356 {
357 let de = RegionDe::deserialize(deserializer)?;
358 let partition_expr = if de.partition_expr.is_empty() {
361 if let Some(LegacyPartition { value_list, .. }) = &de.partition {
362 value_list
363 .first()
364 .map(|expr| String::from_utf8_lossy(expr).to_string())
365 .unwrap_or_default()
366 } else {
367 String::new()
368 }
369 } else {
370 de.partition_expr
371 };
372
373 Ok(Self {
374 id: de.id,
375 name: de.name,
376 attrs: de.attrs,
377 partition_expr,
378 })
379 }
380}
381
382impl Region {
383 #[cfg(any(test, feature = "testing"))]
384 pub fn new_test(id: RegionId) -> Self {
385 Self {
386 id,
387 ..Default::default()
388 }
389 }
390
391 pub fn partition_expr(&self) -> String {
393 self.partition_expr.clone()
394 }
395}
396
397#[derive(Debug, Clone, Deserialize, Serialize, PartialEq)]
398pub struct LegacyPartition {
399 #[serde(serialize_with = "as_utf8_vec", deserialize_with = "from_utf8_vec")]
400 pub column_list: Vec<Vec<u8>>,
401 #[serde(serialize_with = "as_utf8_vec", deserialize_with = "from_utf8_vec")]
402 pub value_list: Vec<Vec<u8>>,
403}
404
405fn as_utf8_vec<S: Serializer>(
406 val: &[Vec<u8>],
407 serializer: S,
408) -> std::result::Result<S::Ok, S::Error> {
409 let mut seq = serializer.serialize_seq(Some(val.len()))?;
410 for v in val {
411 seq.serialize_element(&String::from_utf8_lossy(v))?;
412 }
413 seq.end()
414}
415
416pub fn from_utf8_vec<'de, D>(deserializer: D) -> std::result::Result<Vec<Vec<u8>>, D::Error>
417where
418 D: Deserializer<'de>,
419{
420 let values = Vec::<String>::deserialize(deserializer)?;
421
422 let values = values
423 .into_iter()
424 .map(|value| value.into_bytes())
425 .collect::<Vec<_>>();
426 Ok(values)
427}
428
429#[cfg(test)]
430mod tests {
431 use super::*;
432 use crate::key::RegionRoleSet;
433
434 fn new_test_region_route(region_id: RegionId) -> RegionRoute {
435 RegionRoute {
436 region: Region::new_test(region_id),
437 leader_peer: Some(Peer::new(1, "a1")),
438 follower_peers: vec![Peer::new(2, "a2")],
439 leader_state: None,
440 leader_down_since: None,
441 write_route_policy: None,
442 }
443 }
444
445 #[test]
446 fn test_leader_is_downgraded() {
447 let mut region_route = RegionRoute {
448 region: Region {
449 id: 2.into(),
450 name: "r2".to_string(),
451 attrs: BTreeMap::new(),
452 partition_expr: "".to_string(),
453 },
454 leader_peer: Some(Peer::new(1, "a1")),
455 follower_peers: vec![Peer::new(2, "a2"), Peer::new(3, "a3")],
456 leader_state: None,
457 leader_down_since: None,
458 write_route_policy: None,
459 };
460
461 assert!(!region_route.is_leader_downgrading());
462
463 region_route.downgrade_leader();
464
465 assert!(region_route.is_leader_downgrading());
466 }
467
468 #[test]
469 fn test_region_route_decode() {
470 let region_route = RegionRoute {
471 region: Region {
472 id: 2.into(),
473 name: "r2".to_string(),
474 attrs: BTreeMap::new(),
475 partition_expr: "".to_string(),
476 },
477 leader_peer: Some(Peer::new(1, "a1")),
478 follower_peers: vec![Peer::new(2, "a2"), Peer::new(3, "a3")],
479 leader_state: None,
480 leader_down_since: None,
481 write_route_policy: None,
482 };
483
484 let input = r#"{"region":{"id":2,"name":"r2","partition":null,"attrs":{}},"leader_peer":{"id":1,"addr":"a1"},"follower_peers":[{"id":2,"addr":"a2"},{"id":3,"addr":"a3"}]}"#;
485
486 let decoded: RegionRoute = serde_json::from_str(input).unwrap();
487
488 assert_eq!(decoded, region_route);
489 }
490
491 #[test]
492 fn test_region_route_compatibility() {
493 let region_route = RegionRoute {
494 region: Region {
495 id: 2.into(),
496 name: "r2".to_string(),
497 attrs: BTreeMap::new(),
498 partition_expr: "".to_string(),
499 },
500 leader_peer: Some(Peer::new(1, "a1")),
501 follower_peers: vec![Peer::new(2, "a2"), Peer::new(3, "a3")],
502 leader_state: Some(LeaderState::Downgrading),
503 leader_down_since: None,
504 write_route_policy: None,
505 };
506 let input = r#"{"region":{"id":2,"name":"r2","partition":null,"attrs":{}},"leader_peer":{"id":1,"addr":"a1"},"follower_peers":[{"id":2,"addr":"a2"},{"id":3,"addr":"a3"}],"leader_state":"Downgraded","leader_down_since":null}"#;
507 let decoded: RegionRoute = serde_json::from_str(input).unwrap();
508 assert_eq!(decoded, region_route);
509
510 let region_route = RegionRoute {
511 region: Region {
512 id: 2.into(),
513 name: "r2".to_string(),
514 attrs: BTreeMap::new(),
515 partition_expr: "".to_string(),
516 },
517 leader_peer: Some(Peer::new(1, "a1")),
518 follower_peers: vec![Peer::new(2, "a2"), Peer::new(3, "a3")],
519 leader_state: Some(LeaderState::Downgrading),
520 leader_down_since: None,
521 write_route_policy: None,
522 };
523 let input = r#"{"region":{"id":2,"name":"r2","partition":null,"attrs":{}},"leader_peer":{"id":1,"addr":"a1"},"follower_peers":[{"id":2,"addr":"a2"},{"id":3,"addr":"a3"}],"leader_status":"Downgraded","leader_down_since":null}"#;
524 let decoded: RegionRoute = serde_json::from_str(input).unwrap();
525 assert_eq!(decoded, region_route);
526
527 let region_route = RegionRoute {
528 region: Region {
529 id: 2.into(),
530 name: "r2".to_string(),
531 attrs: BTreeMap::new(),
532 partition_expr: "".to_string(),
533 },
534 leader_peer: Some(Peer::new(1, "a1")),
535 follower_peers: vec![Peer::new(2, "a2"), Peer::new(3, "a3")],
536 leader_state: Some(LeaderState::Downgrading),
537 leader_down_since: None,
538 write_route_policy: None,
539 };
540 let input = r#"{"region":{"id":2,"name":"r2","partition":null,"attrs":{}},"leader_peer":{"id":1,"addr":"a1"},"follower_peers":[{"id":2,"addr":"a2"},{"id":3,"addr":"a3"}],"leader_state":"Downgrading","leader_down_since":null}"#;
541 let decoded: RegionRoute = serde_json::from_str(input).unwrap();
542 assert_eq!(decoded, region_route);
543
544 let region_route = RegionRoute {
545 region: Region {
546 id: 2.into(),
547 name: "r2".to_string(),
548 attrs: BTreeMap::new(),
549 partition_expr: "".to_string(),
550 },
551 leader_peer: Some(Peer::new(1, "a1")),
552 follower_peers: vec![Peer::new(2, "a2"), Peer::new(3, "a3")],
553 leader_state: Some(LeaderState::Downgrading),
554 leader_down_since: None,
555 write_route_policy: None,
556 };
557 let input = r#"{"region":{"id":2,"name":"r2","partition":null,"attrs":{}},"leader_peer":{"id":1,"addr":"a1"},"follower_peers":[{"id":2,"addr":"a2"},{"id":3,"addr":"a3"}],"leader_status":"Downgrading","leader_down_since":null}"#;
558 let decoded: RegionRoute = serde_json::from_str(input).unwrap();
559 assert_eq!(decoded, region_route);
560 }
561
562 #[test]
563 fn test_region_route_write_route_policy_decode_compatibility() {
564 let input = r#"{"region":{"id":2,"name":"r2","partition":null,"attrs":{}},"leader_peer":{"id":1,"addr":"a1"},"follower_peers":[{"id":2,"addr":"a2"}],"write_route_policy":"IgnoreAllWrites"}"#;
565 let decoded: RegionRoute = serde_json::from_str(input).unwrap();
566
567 assert!(decoded.is_ignore_all_writes());
568 }
569
570 #[test]
571 fn test_region_route_write_route_policy_default_not_serialized() {
572 let region_route = RegionRoute {
573 region: Region {
574 id: 2.into(),
575 name: "r2".to_string(),
576 attrs: BTreeMap::new(),
577 partition_expr: "".to_string(),
578 },
579 leader_peer: Some(Peer::new(1, "a1")),
580 follower_peers: vec![Peer::new(2, "a2"), Peer::new(3, "a3")],
581 leader_state: None,
582 leader_down_since: None,
583 write_route_policy: None,
584 };
585
586 let encoded = serde_json::to_string(®ion_route).unwrap();
587 assert!(!encoded.contains("write_route_policy"));
588 }
589
590 #[test]
591 fn test_region_route_write_route_policy_helpers() {
592 let mut region_route = RegionRoute {
593 region: Region::new_test(2.into()),
594 leader_peer: Some(Peer::new(1, "a1")),
595 follower_peers: vec![],
596 leader_state: None,
597 leader_down_since: None,
598 write_route_policy: None,
599 };
600
601 assert!(!region_route.is_ignore_all_writes());
602 region_route.set_ignore_all_writes();
603 assert!(region_route.is_ignore_all_writes());
604 region_route.clear_ignore_all_writes();
605 assert!(!region_route.is_ignore_all_writes());
606 }
607
608 #[test]
609 fn test_leader_region_role_without_leader_peer_returns_none() {
610 let region_route = RegionRoute {
611 leader_peer: None,
612 ..new_test_region_route(RegionId::new(1, 1))
613 };
614
615 assert_eq!(region_route.leader_region_role(), None);
616 }
617
618 #[test]
619 fn test_leader_region_role_variants() {
620 let normal = new_test_region_route(RegionId::new(1, 1));
621 let mut downgrading = new_test_region_route(RegionId::new(1, 2));
622 downgrading.leader_state = Some(LeaderState::Downgrading);
623 let mut staging = new_test_region_route(RegionId::new(1, 3));
624 staging.leader_state = Some(LeaderState::Staging);
625
626 assert_eq!(normal.leader_region_role(), Some(RegionRole::Leader));
627 assert_eq!(
628 downgrading.leader_region_role(),
629 Some(RegionRole::DowngradingLeader)
630 );
631 assert_eq!(
632 staging.leader_region_role(),
633 Some(RegionRole::StagingLeader)
634 );
635 }
636
637 #[test]
638 fn test_operating_leader_region_roles_returns_expected_roles() {
639 let no_leader_region = RegionRoute {
640 leader_peer: None,
641 ..new_test_region_route(RegionId::new(1, 4))
642 };
643 let mut downgrading = new_test_region_route(RegionId::new(1, 2));
644 downgrading.leader_peer = Some(Peer::new(2, "a2"));
645 downgrading.leader_state = Some(LeaderState::Downgrading);
646 let mut staging = new_test_region_route(RegionId::new(1, 3));
647 staging.leader_peer = Some(Peer::new(3, "a3"));
648 staging.leader_state = Some(LeaderState::Staging);
649
650 let roles = operating_leader_region_roles(&[
651 new_test_region_route(RegionId::new(1, 1)),
652 downgrading,
653 staging,
654 no_leader_region,
655 ]);
656
657 assert_eq!(
658 roles,
659 vec![
660 (RegionId::new(1, 1), 1, RegionRole::Leader),
661 (RegionId::new(1, 2), 2, RegionRole::DowngradingLeader),
662 (RegionId::new(1, 3), 3, RegionRole::StagingLeader),
663 ]
664 );
665 }
666
667 #[test]
668 fn test_region_distribution() {
669 let region_routes = vec![
670 RegionRoute {
671 region: Region {
672 id: RegionId::new(1, 1),
673 name: "r1".to_string(),
674 attrs: BTreeMap::new(),
675 partition_expr: "".to_string(),
676 },
677 leader_peer: Some(Peer::new(1, "a1")),
678 follower_peers: vec![Peer::new(2, "a2"), Peer::new(3, "a3")],
679 leader_state: None,
680 leader_down_since: None,
681 write_route_policy: None,
682 },
683 RegionRoute {
684 region: Region {
685 id: RegionId::new(1, 2),
686 name: "r2".to_string(),
687 attrs: BTreeMap::new(),
688 partition_expr: "".to_string(),
689 },
690 leader_peer: Some(Peer::new(2, "a2")),
691 follower_peers: vec![Peer::new(1, "a1"), Peer::new(3, "a3")],
692 leader_state: None,
693 leader_down_since: None,
694 write_route_policy: None,
695 },
696 ];
697
698 let distribution = region_distribution(®ion_routes);
699 assert_eq!(distribution.len(), 3);
700 assert_eq!(distribution[&1], RegionRoleSet::new(vec![1], vec![2]));
701 assert_eq!(distribution[&2], RegionRoleSet::new(vec![2], vec![1]));
702 assert_eq!(distribution[&3], RegionRoleSet::new(vec![], vec![1, 2]));
703 }
704
705 #[test]
706 fn test_de_serialize_partition() {
707 let p = LegacyPartition {
708 column_list: vec![b"a".to_vec(), b"b".to_vec()],
709 value_list: vec![b"hi".to_vec(), b",".to_vec()],
710 };
711
712 let output = serde_json::to_string(&p).unwrap();
713 let got: LegacyPartition = serde_json::from_str(&output).unwrap();
714
715 assert_eq!(got, p);
716 }
717}