1use std::collections::HashMap;
16use std::sync::Arc;
17
18use api::v1::meta::{HeartbeatRequest, RegionLease, Role};
19use async_trait::async_trait;
20use common_meta::key::TableMetadataManagerRef;
21use common_meta::region_keeper::MemoryRegionKeeperRef;
22use common_telemetry::error;
23use store_api::region_engine::GrantedRegion;
24use store_api::storage::RegionId;
25
26use crate::error::Result;
27use crate::handler::{HandleControl, HeartbeatAccumulator, HeartbeatHandler};
28use crate::metasrv::Context;
29use crate::region::RegionLeaseKeeper;
30use crate::region::lease_keeper::{
31 RegionLeaseInfo, RegionLeaseKeeperRef, RenewRegionLeasesResponse,
32};
33
34pub struct RegionLeaseHandler {
35 region_lease_seconds: u64,
36 region_lease_keeper: RegionLeaseKeeperRef,
37 customized_region_lease_renewer: Option<CustomizedRegionLeaseRenewerRef>,
38}
39
40pub type CustomizedRegionLeaseRenewerRef = Arc<dyn CustomizedRegionLeaseRenewer>;
41
42pub trait CustomizedRegionLeaseRenewer: Send + Sync {
43 fn renew(
44 &self,
45 ctx: &mut Context,
46 regions: HashMap<RegionId, RegionLeaseInfo>,
47 ) -> Vec<GrantedRegion>;
48}
49
50impl RegionLeaseHandler {
51 pub fn new(
52 region_lease_seconds: u64,
53 table_metadata_manager: TableMetadataManagerRef,
54 memory_region_keeper: MemoryRegionKeeperRef,
55 customized_region_lease_renewer: Option<CustomizedRegionLeaseRenewerRef>,
56 ) -> Self {
57 let region_lease_keeper =
58 RegionLeaseKeeper::new(table_metadata_manager, memory_region_keeper.clone());
59
60 Self {
61 region_lease_seconds,
62 region_lease_keeper: Arc::new(region_lease_keeper),
63 customized_region_lease_renewer,
64 }
65 }
66}
67
68#[async_trait]
69impl HeartbeatHandler for RegionLeaseHandler {
70 fn is_acceptable(&self, role: Role) -> bool {
71 role == Role::Datanode
72 }
73
74 async fn handle(
75 &self,
76 req: &HeartbeatRequest,
77 ctx: &mut Context,
78 acc: &mut HeartbeatAccumulator,
79 ) -> Result<HandleControl> {
80 let Some(stat) = acc.stat.as_ref() else {
81 return Ok(HandleControl::Continue);
82 };
83
84 let regions = stat.regions();
85 let datanode_id = stat.id;
86
87 match self
88 .region_lease_keeper
89 .renew_region_leases(datanode_id, ®ions)
90 .await
91 {
92 Ok(RenewRegionLeasesResponse {
93 non_exists,
94 renewed,
95 }) => {
96 let renewed = if let Some(renewer) = &self.customized_region_lease_renewer {
97 renewer
98 .renew(ctx, renewed)
99 .into_iter()
100 .map(|region| region.into())
101 .collect()
102 } else {
103 renewed
104 .into_iter()
105 .map(|(region_id, region_lease_info)| {
106 GrantedRegion::new(region_id, region_lease_info.role).into()
107 })
108 .collect::<Vec<_>>()
109 };
110
111 acc.region_lease = Some(RegionLease {
112 regions: renewed,
113 duration_since_epoch: req.duration_since_epoch,
114 lease_seconds: self.region_lease_seconds,
115 closeable_region_ids: non_exists.iter().map(|region| region.as_u64()).collect(),
116 });
117 acc.inactive_region_ids = non_exists;
118 }
119 Err(e) => {
120 error!(e; "Failed to renew region leases for datanode: {datanode_id:?}, regions: {:?}", regions);
121 }
124 }
125
126 Ok(HandleControl::Continue)
127 }
128}
129
130#[cfg(test)]
131mod test {
132
133 use std::collections::{HashMap, HashSet};
134 use std::sync::Arc;
135
136 use common_meta::datanode::{RegionManifestInfo, RegionStat, Stat};
137 use common_meta::distributed_time_constants::default_distributed_time_constants;
138 use common_meta::key::TableMetadataManager;
139 use common_meta::key::table_route::TableRouteValue;
140 use common_meta::key::test_utils::new_test_table_info;
141 use common_meta::kv_backend::memory::MemoryKvBackend;
142 use common_meta::kv_backend::test_util::MockKvBackendBuilder;
143 use common_meta::peer::Peer;
144 use common_meta::region_keeper::MemoryRegionKeeper;
145 use common_meta::rpc::router::{LeaderState, Region, RegionRoute};
146 use store_api::region_engine::RegionRole;
147 use store_api::storage::RegionId;
148
149 use super::*;
150 use crate::metasrv::builder::MetasrvBuilder;
151
152 fn new_test_keeper() -> RegionLeaseKeeper {
153 let store = Arc::new(MemoryKvBackend::new());
154
155 let table_metadata_manager = Arc::new(TableMetadataManager::new(store));
156
157 let memory_region_keeper = Arc::new(MemoryRegionKeeper::default());
158 RegionLeaseKeeper::new(table_metadata_manager, memory_region_keeper)
159 }
160
161 fn new_empty_region_stat(region_id: RegionId, role: RegionRole) -> RegionStat {
162 RegionStat {
163 id: region_id,
164 role,
165 rcus: 0,
166 wcus: 0,
167 approximate_bytes: 0,
168 engine: String::new(),
169 num_rows: 0,
170 memtable_size: 0,
171 manifest_size: 0,
172 sst_size: 0,
173 sst_num: 0,
174 index_size: 0,
175 region_manifest: RegionManifestInfo::Mito {
176 manifest_version: 0,
177 flushed_entry_id: 0,
178 file_removed_cnt: 0,
179 },
180 data_topic_latest_entry_id: 0,
181 metadata_topic_latest_entry_id: 0,
182 written_bytes: 0,
183 query_cpu_time: 0,
184 query_scanned_bytes: 0,
185 min_timestamp: None,
186 max_timestamp: None,
187 }
188 }
189
190 #[tokio::test]
191 async fn test_handle_upgradable_follower() {
192 let datanode_id = 1;
193 let region_number = 1u32;
194 let table_id = 10;
195 let region_id = RegionId::new(table_id, region_number);
196 let another_region_id = RegionId::new(table_id, region_number + 1);
197 let peer = Peer::empty(datanode_id);
198 let follower_peer = Peer::empty(datanode_id + 1);
199 let table_info = new_test_table_info(table_id);
200
201 let region_routes = vec![RegionRoute {
202 region: Region::new_test(region_id),
203 leader_peer: Some(peer.clone()),
204 follower_peers: vec![follower_peer.clone()],
205 ..Default::default()
206 }];
207
208 let keeper = new_test_keeper();
209 let table_metadata_manager = keeper.table_metadata_manager();
210
211 table_metadata_manager
212 .create_table_metadata(
213 table_info,
214 TableRouteValue::physical(region_routes),
215 HashMap::default(),
216 )
217 .await
218 .unwrap();
219
220 let builder = MetasrvBuilder::new();
221 let metasrv = builder.build().await.unwrap();
222 let ctx = &mut metasrv.new_ctx();
223
224 let acc = &mut HeartbeatAccumulator::default();
225
226 acc.stat = Some(Stat {
227 id: peer.id,
228 region_stats: vec![
229 new_empty_region_stat(region_id, RegionRole::Follower),
230 new_empty_region_stat(another_region_id, RegionRole::Follower),
231 ],
232 ..Default::default()
233 });
234
235 let req = HeartbeatRequest {
236 duration_since_epoch: 1234,
237 ..Default::default()
238 };
239
240 let opening_region_keeper = Arc::new(MemoryRegionKeeper::default());
241
242 let handler = RegionLeaseHandler::new(
243 default_distributed_time_constants().region_lease.as_secs(),
244 table_metadata_manager.clone(),
245 opening_region_keeper.clone(),
246 None,
247 );
248
249 handler.handle(&req, ctx, acc).await.unwrap();
250
251 assert_region_lease(acc, vec![GrantedRegion::new(region_id, RegionRole::Leader)]);
252 assert_eq!(acc.inactive_region_ids, HashSet::from([another_region_id]));
253 assert_eq!(
254 acc.region_lease.as_ref().unwrap().closeable_region_ids,
255 vec![another_region_id]
256 );
257
258 let acc = &mut HeartbeatAccumulator::default();
259
260 acc.stat = Some(Stat {
261 id: follower_peer.id,
262 region_stats: vec![
263 new_empty_region_stat(region_id, RegionRole::Follower),
264 new_empty_region_stat(another_region_id, RegionRole::Follower),
265 ],
266 ..Default::default()
267 });
268
269 handler.handle(&req, ctx, acc).await.unwrap();
270
271 assert_eq!(
272 acc.region_lease.as_ref().unwrap().lease_seconds,
273 default_distributed_time_constants().region_lease.as_secs()
274 );
275
276 assert_region_lease(
277 acc,
278 vec![GrantedRegion::new(region_id, RegionRole::Follower)],
279 );
280 assert_eq!(acc.inactive_region_ids, HashSet::from([another_region_id]));
281 assert_eq!(
282 acc.region_lease.as_ref().unwrap().closeable_region_ids,
283 vec![another_region_id]
284 );
285
286 let opening_region_id = RegionId::new(table_id, region_number + 2);
287 let _guard = opening_region_keeper
288 .register_with_role(follower_peer.id, opening_region_id, RegionRole::Follower)
289 .unwrap();
290
291 let acc = &mut HeartbeatAccumulator::default();
292
293 acc.stat = Some(Stat {
294 id: follower_peer.id,
295 region_stats: vec![
296 new_empty_region_stat(region_id, RegionRole::Follower),
297 new_empty_region_stat(another_region_id, RegionRole::Follower),
298 new_empty_region_stat(opening_region_id, RegionRole::Follower),
299 ],
300 ..Default::default()
301 });
302
303 handler.handle(&req, ctx, acc).await.unwrap();
304
305 assert_eq!(
306 acc.region_lease.as_ref().unwrap().lease_seconds,
307 default_distributed_time_constants().region_lease.as_secs()
308 );
309
310 assert_region_lease(
311 acc,
312 vec![
313 GrantedRegion::new(region_id, RegionRole::Follower),
314 GrantedRegion::new(opening_region_id, RegionRole::Follower),
315 ],
316 );
317 assert_eq!(acc.inactive_region_ids, HashSet::from([another_region_id]));
318 assert_eq!(
319 acc.region_lease.as_ref().unwrap().closeable_region_ids,
320 vec![another_region_id]
321 );
322 }
323
324 #[tokio::test]
325
326 async fn test_handle_downgradable_leader() {
327 let datanode_id = 1;
328 let region_number = 1u32;
329 let table_id = 10;
330 let region_id = RegionId::new(table_id, region_number);
331 let another_region_id = RegionId::new(table_id, region_number + 1);
332 let no_exist_region_id = RegionId::new(table_id, region_number + 2);
333 let peer = Peer::empty(datanode_id);
334 let follower_peer = Peer::empty(datanode_id + 1);
335 let table_info = new_test_table_info(table_id);
336
337 let region_routes = vec![
338 RegionRoute {
339 region: Region::new_test(region_id),
340 leader_peer: Some(peer.clone()),
341 follower_peers: vec![follower_peer.clone()],
342 leader_state: Some(LeaderState::Downgrading),
343 leader_down_since: Some(1),
344 write_route_policy: None,
345 },
346 RegionRoute {
347 region: Region::new_test(another_region_id),
348 leader_peer: Some(peer.clone()),
349 ..Default::default()
350 },
351 ];
352
353 let keeper = new_test_keeper();
354 let table_metadata_manager = keeper.table_metadata_manager();
355
356 table_metadata_manager
357 .create_table_metadata(
358 table_info,
359 TableRouteValue::physical(region_routes),
360 HashMap::default(),
361 )
362 .await
363 .unwrap();
364
365 let builder = MetasrvBuilder::new();
366 let metasrv = builder.build().await.unwrap();
367 let ctx = &mut metasrv.new_ctx();
368
369 let req = HeartbeatRequest {
370 duration_since_epoch: 1234,
371 ..Default::default()
372 };
373
374 let acc = &mut HeartbeatAccumulator::default();
375
376 acc.stat = Some(Stat {
377 id: peer.id,
378 region_stats: vec![
379 new_empty_region_stat(region_id, RegionRole::Leader),
380 new_empty_region_stat(another_region_id, RegionRole::Leader),
381 new_empty_region_stat(no_exist_region_id, RegionRole::Leader),
382 ],
383 ..Default::default()
384 });
385
386 let handler = RegionLeaseHandler::new(
387 default_distributed_time_constants().region_lease.as_secs(),
388 table_metadata_manager.clone(),
389 Default::default(),
390 None,
391 );
392
393 handler.handle(&req, ctx, acc).await.unwrap();
394
395 assert_region_lease(
396 acc,
397 vec![
398 GrantedRegion::new(region_id, RegionRole::DowngradingLeader),
399 GrantedRegion::new(another_region_id, RegionRole::Leader),
400 ],
401 );
402 assert_eq!(acc.inactive_region_ids, HashSet::from([no_exist_region_id]));
403 }
404
405 #[tokio::test]
406 async fn test_handle_staging_leader() {
407 let datanode_id = 1;
408 let region_number = 1u32;
409 let table_id = 10;
410 let region_id = RegionId::new(table_id, region_number);
411 let peer = Peer::empty(datanode_id);
412 let table_info = new_test_table_info(table_id);
413
414 let region_routes = vec![RegionRoute {
415 region: Region::new_test(region_id),
416 leader_peer: Some(peer.clone()),
417 leader_state: Some(LeaderState::Staging),
418 ..Default::default()
419 }];
420
421 let keeper = new_test_keeper();
422 let table_metadata_manager = keeper.table_metadata_manager();
423
424 table_metadata_manager
425 .create_table_metadata(
426 table_info,
427 TableRouteValue::physical(region_routes),
428 HashMap::default(),
429 )
430 .await
431 .unwrap();
432
433 let builder = MetasrvBuilder::new();
434 let metasrv = builder.build().await.unwrap();
435 let ctx = &mut metasrv.new_ctx();
436
437 let req = HeartbeatRequest {
438 duration_since_epoch: 1234,
439 ..Default::default()
440 };
441
442 let acc = &mut HeartbeatAccumulator::default();
443 acc.stat = Some(Stat {
444 id: peer.id,
445 region_stats: vec![new_empty_region_stat(region_id, RegionRole::StagingLeader)],
446 ..Default::default()
447 });
448
449 let handler = RegionLeaseHandler::new(
450 default_distributed_time_constants().region_lease.as_secs(),
451 table_metadata_manager.clone(),
452 Default::default(),
453 None,
454 );
455
456 handler.handle(&req, ctx, acc).await.unwrap();
457
458 assert_region_lease(
459 acc,
460 vec![GrantedRegion::new(region_id, RegionRole::StagingLeader)],
461 );
462 }
463
464 fn assert_region_lease(acc: &HeartbeatAccumulator, expected: Vec<GrantedRegion>) {
465 let region_lease = acc.region_lease.as_ref().unwrap().clone();
466 let granted: Vec<GrantedRegion> = region_lease
467 .regions
468 .into_iter()
469 .map(Into::into)
470 .collect::<Vec<_>>();
471
472 let granted = granted
473 .into_iter()
474 .map(|region| (region.region_id, region))
475 .collect::<HashMap<_, _>>();
476
477 let expected = expected
478 .into_iter()
479 .map(|region| (region.region_id, region))
480 .collect::<HashMap<_, _>>();
481
482 assert_eq!(granted, expected);
483 }
484
485 #[tokio::test]
486 async fn test_handle_renew_region_lease_failure() {
487 common_telemetry::init_default_ut_logging();
488 let kv = MockKvBackendBuilder::default()
489 .batch_get_fn(Arc::new(|_| {
490 common_meta::error::UnexpectedSnafu {
491 err_msg: "mock err",
492 }
493 .fail()
494 }) as _)
495 .build()
496 .unwrap();
497 let kvbackend = Arc::new(kv);
498 let table_metadata_manager = Arc::new(TableMetadataManager::new(kvbackend));
499
500 let datanode_id = 1;
501 let region_number = 1u32;
502 let table_id = 10;
503 let region_id = RegionId::new(table_id, region_number);
504 let another_region_id = RegionId::new(table_id, region_number + 1);
505 let no_exist_region_id = RegionId::new(table_id, region_number + 2);
506 let peer = Peer::empty(datanode_id);
507
508 let builder = MetasrvBuilder::new();
509 let metasrv = builder.build().await.unwrap();
510 let ctx = &mut metasrv.new_ctx();
511
512 let req = HeartbeatRequest {
513 duration_since_epoch: 1234,
514 ..Default::default()
515 };
516
517 let acc = &mut HeartbeatAccumulator::default();
518 acc.stat = Some(Stat {
519 id: peer.id,
520 region_stats: vec![
521 new_empty_region_stat(region_id, RegionRole::Leader),
522 new_empty_region_stat(another_region_id, RegionRole::Leader),
523 new_empty_region_stat(no_exist_region_id, RegionRole::Leader),
524 ],
525 ..Default::default()
526 });
527 let handler = RegionLeaseHandler::new(
528 default_distributed_time_constants().region_lease.as_secs(),
529 table_metadata_manager.clone(),
530 Default::default(),
531 None,
532 );
533 handler.handle(&req, ctx, acc).await.unwrap();
534
535 assert!(acc.region_lease.is_none());
536 assert!(acc.inactive_region_ids.is_empty());
537 }
538}