1use api::v1::meta::{HeartbeatRequest, Role};
16use common_meta::region_registry::{LeaderRegion, LeaderRegionManifestInfo};
17use store_api::region_engine::RegionRole;
18
19use crate::error::Result;
20use crate::handler::{HandleControl, HeartbeatAccumulator, HeartbeatHandler};
21use crate::metasrv::Context;
22
23pub struct CollectLeaderRegionHandler;
24
25#[async_trait::async_trait]
26impl HeartbeatHandler for CollectLeaderRegionHandler {
27 fn is_acceptable(&self, role: Role) -> bool {
28 role == Role::Datanode
29 }
30
31 async fn handle(
32 &self,
33 _req: &HeartbeatRequest,
34 ctx: &mut Context,
35 acc: &mut HeartbeatAccumulator,
36 ) -> Result<HandleControl> {
37 let Some(current_stat) = acc.stat.as_ref() else {
38 return Ok(HandleControl::Continue);
39 };
40
41 let mut key_values = Vec::with_capacity(current_stat.region_stats.len());
42 for stat in current_stat.region_stats.iter() {
43 if !matches!(stat.role, RegionRole::Leader | RegionRole::StagingLeader) {
44 continue;
45 }
46
47 let manifest = LeaderRegionManifestInfo::from_region_stat(stat);
48 let value = LeaderRegion {
49 datanode_id: current_stat.id,
50 manifest,
51 };
52 key_values.push((stat.id, value));
53 }
54 ctx.leader_region_registry.batch_put(key_values);
55
56 Ok(HandleControl::Continue)
57 }
58}
59
60#[cfg(test)]
61mod tests {
62 use common_meta::datanode::{RegionManifestInfo, RegionStat, Stat};
63 use common_meta::region_registry::LeaderRegionManifestInfo;
64 use store_api::region_engine::RegionRole;
65 use store_api::storage::RegionId;
66
67 use super::*;
68 use crate::handler::test_utils::TestEnv;
69
70 fn new_region_stat(id: RegionId, manifest_version: u64, role: RegionRole) -> RegionStat {
71 RegionStat {
72 id,
73 region_manifest: RegionManifestInfo::Mito {
74 manifest_version,
75 flushed_entry_id: 0,
76 file_removed_cnt: 0,
77 },
78 rcus: 0,
79 wcus: 0,
80 approximate_bytes: 0,
81 engine: "mito".to_string(),
82 role,
83 num_rows: 0,
84 memtable_size: 0,
85 manifest_size: 0,
86 sst_size: 0,
87 sst_num: 0,
88 index_size: 0,
89 data_topic_latest_entry_id: 0,
90 metadata_topic_latest_entry_id: 0,
91 written_bytes: 0,
92 query_cpu_time: 0,
93 query_scanned_bytes: 0,
94 }
95 }
96
97 #[tokio::test]
98 async fn test_handle_collect_leader_region() {
99 let env = TestEnv::new();
100 let mut ctx = env.ctx();
101
102 let mut acc = HeartbeatAccumulator {
103 stat: Some(Stat {
104 id: 1,
105 region_stats: vec![
106 new_region_stat(RegionId::new(1, 1), 1, RegionRole::Leader),
107 new_region_stat(RegionId::new(1, 2), 2, RegionRole::Follower),
108 ],
109 addr: "127.0.0.1:0000".to_string(),
110 region_num: 2,
111 ..Default::default()
112 }),
113 ..Default::default()
114 };
115
116 let handler = CollectLeaderRegionHandler;
117 let control = handler
118 .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
119 .await
120 .unwrap();
121
122 assert_eq!(control, HandleControl::Continue);
123 let regions = ctx
124 .leader_region_registry
125 .batch_get(vec![RegionId::new(1, 1), RegionId::new(1, 2)].into_iter());
126 assert_eq!(regions.len(), 1);
127 assert_eq!(
128 regions.get(&RegionId::new(1, 1)),
129 Some(&LeaderRegion {
130 datanode_id: 1,
131 manifest: LeaderRegionManifestInfo::Mito {
132 manifest_version: 1,
133 flushed_entry_id: 0,
134 topic_latest_entry_id: 0,
135 },
136 })
137 );
138
139 acc.stat = Some(Stat {
141 id: 1,
142 region_stats: vec![new_region_stat(RegionId::new(1, 1), 2, RegionRole::Leader)],
143 timestamp_millis: 0,
144 addr: "127.0.0.1:0000".to_string(),
145 region_num: 1,
146 node_epoch: 0,
147 ..Default::default()
148 });
149 let control = handler
150 .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
151 .await
152 .unwrap();
153
154 assert_eq!(control, HandleControl::Continue);
155 let regions = ctx
156 .leader_region_registry
157 .batch_get(vec![RegionId::new(1, 1)].into_iter());
158 assert_eq!(regions.len(), 1);
159 assert_eq!(
160 regions.get(&RegionId::new(1, 1)),
161 Some(&LeaderRegion {
162 datanode_id: 1,
163 manifest: LeaderRegionManifestInfo::Mito {
164 manifest_version: 2,
165 flushed_entry_id: 0,
166 topic_latest_entry_id: 0,
167 },
168 })
169 );
170
171 acc.stat = Some(Stat {
173 id: 1,
174 region_stats: vec![new_region_stat(RegionId::new(1, 1), 1, RegionRole::Leader)],
175 timestamp_millis: 0,
176 addr: "127.0.0.1:0000".to_string(),
177 region_num: 1,
178 node_epoch: 0,
179 ..Default::default()
180 });
181 let control = handler
182 .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
183 .await
184 .unwrap();
185
186 assert_eq!(control, HandleControl::Continue);
187 let regions = ctx
188 .leader_region_registry
189 .batch_get(vec![RegionId::new(1, 1)].into_iter());
190 assert_eq!(regions.len(), 1);
191 assert_eq!(
192 regions.get(&RegionId::new(1, 1)),
193 Some(&LeaderRegion {
195 datanode_id: 1,
196 manifest: LeaderRegionManifestInfo::Mito {
197 manifest_version: 2,
198 flushed_entry_id: 0,
199 topic_latest_entry_id: 0,
200 },
201 })
202 );
203 }
204}