1use std::any::Any;
16use std::collections::HashSet;
17use std::time::Duration;
18
19use api::v1::meta::MailboxMessage;
20use common_meta::instruction::{
21 Instruction, InstructionReply, UpgradeRegion, UpgradeRegionReply, UpgradeRegionsReply,
22};
23use common_meta::lock_key::RemoteWalLock;
24use common_meta::wal_provider::extract_topic_from_wal_options;
25use common_procedure::{Context as ProcedureContext, Status};
26use common_telemetry::tracing_context::TracingContext;
27use common_telemetry::{error, info};
28use common_wal::options::WalOptions;
29use serde::{Deserialize, Serialize};
30use snafu::{OptionExt, ResultExt, ensure};
31use store_api::metric_engine_consts::METRIC_ENGINE_NAME;
32use tokio::time::{Instant, sleep};
33
34use crate::error::{self, Result};
35use crate::handler::HeartbeatMailbox;
36use crate::procedure::region_migration::update_metadata::UpdateMetadata;
37use crate::procedure::region_migration::{Context, State};
38use crate::procedure::utils::instruction_error_result;
39use crate::service::mailbox::Channel;
40
41#[derive(Debug, Serialize, Deserialize)]
42pub struct UpgradeCandidateRegion {
43 pub(crate) optimistic_retry: usize,
45 pub(crate) retry_initial_interval: Duration,
47 pub(crate) require_ready: bool,
50}
51
52impl Default for UpgradeCandidateRegion {
53 fn default() -> Self {
54 Self {
55 optimistic_retry: 3,
56 retry_initial_interval: Duration::from_millis(500),
57 require_ready: true,
58 }
59 }
60}
61
62#[async_trait::async_trait]
63#[typetag::serde]
64impl State for UpgradeCandidateRegion {
65 async fn next(
66 &mut self,
67 ctx: &mut Context,
68 procedure_ctx: &ProcedureContext,
69 ) -> Result<(Box<dyn State>, Status)> {
70 let now = Instant::now();
71
72 let topics = self.get_kafka_topics(ctx).await?;
73 if self
74 .upgrade_region_with_retry(ctx, procedure_ctx, topics)
75 .await
76 {
77 ctx.update_upgrade_candidate_region_elapsed(now);
78 Ok((Box::new(UpdateMetadata::Upgrade), Status::executing(false)))
79 } else {
80 ctx.update_upgrade_candidate_region_elapsed(now);
81 Ok((Box::new(UpdateMetadata::Rollback), Status::executing(false)))
82 }
83 }
84
85 fn as_any(&self) -> &dyn Any {
86 self
87 }
88}
89
90impl UpgradeCandidateRegion {
91 async fn get_kafka_topics(&self, ctx: &mut Context) -> Result<HashSet<String>> {
92 let table_regions = ctx.persistent_ctx.table_regions();
93 let datanode_table_values = ctx.get_from_peer_datanode_table_values().await?;
94 let mut topics = HashSet::new();
95 for (table_id, regions) in table_regions {
96 let Some(datanode_table_value) = datanode_table_values.get(&table_id) else {
97 continue;
98 };
99
100 let region_wal_options = &datanode_table_value.region_info.region_wal_options;
101
102 for region_id in regions {
103 let Some(WalOptions::Kafka(kafka_wal_options)) =
104 region_wal_options.get(®ion_id.region_number())
105 else {
106 continue;
107 };
108 if !topics.contains(&kafka_wal_options.topic) {
109 topics.insert(kafka_wal_options.topic.clone());
110 }
111 }
112 }
113
114 Ok(topics)
115 }
116
117 async fn build_upgrade_region_instruction(
119 &self,
120 ctx: &mut Context,
121 replay_timeout: Duration,
122 ) -> Result<Instruction> {
123 let region_ids = ctx.persistent_ctx.region_ids.clone();
124 let datanode_table_values = ctx.get_from_peer_datanode_table_values().await?;
125 let mut region_topic = Vec::with_capacity(region_ids.len());
126 for region_id in region_ids.iter() {
127 let table_id = region_id.table_id();
128 if let Some(datanode_table_value) = datanode_table_values.get(&table_id)
129 && let Some(topic) = extract_topic_from_wal_options(
130 *region_id,
131 &datanode_table_value.region_info.region_wal_options,
132 )
133 {
134 let is_metric_engine =
135 datanode_table_value.region_info.engine == METRIC_ENGINE_NAME;
136 region_topic.push((*region_id, topic, is_metric_engine));
137 }
138 }
139
140 let replay_checkpoints = ctx
141 .get_replay_checkpoints_with_topic_pruned_entry_ids(®ion_topic)
142 .await?;
143 let mut upgrade_regions = Vec::with_capacity(region_ids.len());
145 for region_id in region_ids {
146 let last_entry_id = ctx
147 .volatile_ctx
148 .leader_region_last_entry_ids
149 .get(®ion_id)
150 .copied();
151 let metadata_last_entry_id = ctx
152 .volatile_ctx
153 .leader_region_metadata_last_entry_ids
154 .get(®ion_id)
155 .copied();
156 let checkpoint = replay_checkpoints.get(®ion_id).copied();
157 upgrade_regions.push(UpgradeRegion {
158 region_id,
159 last_entry_id,
160 metadata_last_entry_id,
161 replay_timeout,
162 location_id: Some(ctx.persistent_ctx.from_peer.id),
163 replay_entry_id: checkpoint.map(|c| c.entry_id),
164 metadata_replay_entry_id: checkpoint.and_then(|c| c.metadata_entry_id),
165 });
166 }
167
168 Ok(Instruction::UpgradeRegions(upgrade_regions))
169 }
170
171 fn handle_upgrade_region_reply(
172 &self,
173 ctx: &mut Context,
174 UpgradeRegionReply {
175 region_id,
176 ready,
177 exists,
178 error,
179 }: &UpgradeRegionReply,
180 now: &Instant,
181 ) -> Result<()> {
182 let candidate = &ctx.persistent_ctx.to_peer;
183 if let Some(error) = error {
184 return instruction_error_result(
185 error,
186 format!(
187 "Failed to upgrade the region {} on datanode {:?}, error: {:?}, elapsed: {:?}",
188 region_id,
189 candidate,
190 error,
191 now.elapsed()
192 ),
193 );
194 }
195
196 ensure!(
197 exists,
198 error::UnexpectedSnafu {
199 violated: format!(
200 "Candidate region {} doesn't exist on datanode {:?}",
201 region_id, candidate
202 )
203 }
204 );
205
206 if self.require_ready && !ready {
207 return error::RetryLaterSnafu {
208 reason: format!(
209 "Candidate region {} still replaying the wal on datanode {:?}, elapsed: {:?}",
210 region_id,
211 candidate,
212 now.elapsed()
213 ),
214 }
215 .fail();
216 }
217
218 Ok(())
219 }
220
221 async fn upgrade_region(&self, ctx: &mut Context) -> Result<()> {
236 let operation_timeout =
237 ctx.next_operation_timeout()
238 .context(error::ExceededDeadlineSnafu {
239 operation: "Upgrade region",
240 })?;
241 let upgrade_instruction = self
242 .build_upgrade_region_instruction(ctx, operation_timeout)
243 .await?;
244
245 let pc = &ctx.persistent_ctx;
246 let region_ids = &pc.region_ids;
247 let candidate = &pc.to_peer;
248
249 let tracing_ctx = TracingContext::from_current_span();
250 let msg = MailboxMessage::json_message(
251 &format!("Upgrade candidate regions: {:?}", region_ids),
252 &format!("Metasrv@{}", ctx.server_addr()),
253 &format!("Datanode-{}@{}", candidate.id, candidate.addr),
254 common_time::util::current_time_millis(),
255 &upgrade_instruction,
256 Some(tracing_ctx.to_w3c()),
257 )
258 .with_context(|_| error::SerializeToJsonSnafu {
259 input: upgrade_instruction.to_string(),
260 })?;
261
262 let ch = Channel::Datanode(candidate.id);
263 let receiver = ctx.mailbox.send(&ch, msg, operation_timeout).await?;
264
265 let now = Instant::now();
266 match receiver.await {
267 Ok(msg) => {
268 let reply = HeartbeatMailbox::json_reply(&msg)?;
269 info!(
270 "Received upgrade region reply: {:?}, regions: {:?}, elapsed: {:?}",
271 reply,
272 region_ids,
273 now.elapsed()
274 );
275 let InstructionReply::UpgradeRegions(UpgradeRegionsReply { replies }) = reply
276 else {
277 return error::UnexpectedInstructionReplySnafu {
278 mailbox_message: msg.to_string(),
279 reason: "Unexpected reply of the upgrade region instruction",
280 }
281 .fail();
282 };
283 for reply in replies {
284 self.handle_upgrade_region_reply(ctx, &reply, &now)?;
285 }
286 Ok(())
287 }
288 Err(error::Error::MailboxTimeout { .. }) => {
289 let reason = format!(
290 "Mailbox received timeout for upgrade candidate regions {region_ids:?} on datanode {:?}, elapsed: {:?}",
291 candidate,
292 now.elapsed()
293 );
294 error::RetryLaterSnafu { reason }.fail()
295 }
296 Err(err) => Err(err),
297 }
298 }
299
300 async fn upgrade_region_with_retry(
304 &self,
305 ctx: &mut Context,
306 procedure_ctx: &ProcedureContext,
307 topics: HashSet<String>,
308 ) -> bool {
309 let mut retry = 0;
310 let mut upgraded = false;
311
312 let mut guards = Vec::with_capacity(topics.len());
313 loop {
314 let timer = Instant::now();
315 for topic in &topics {
317 guards.push(
318 procedure_ctx
319 .provider
320 .acquire_lock(&(RemoteWalLock::Read(topic.clone()).into()))
321 .await,
322 );
323 }
324
325 if let Err(err) = self.upgrade_region(ctx).await {
326 retry += 1;
327 ctx.update_operations_elapsed(timer);
328 if matches!(err, error::Error::ExceededDeadline { .. }) {
329 error!("Failed to upgrade region, exceeded deadline");
330 break;
331 } else if err.is_retryable() && retry < self.optimistic_retry {
332 error!("Failed to upgrade region, error: {err:?}, retry later");
333 sleep(self.retry_initial_interval).await;
334 } else {
335 error!("Failed to upgrade region, error: {err:?}");
336 break;
337 }
338 } else {
339 ctx.update_operations_elapsed(timer);
340 upgraded = true;
341 break;
342 }
343 }
344
345 upgraded
346 }
347}
348
349#[cfg(test)]
350mod tests {
351 use std::assert_matches;
352 use std::collections::HashMap;
353
354 use common_meta::key::table_route::TableRouteValue;
355 use common_meta::key::test_utils::new_test_table_info;
356 use common_meta::key::topic_name::TopicNameKey;
357 use common_meta::key::topic_region::{ReplayCheckpoint, TopicRegionKey, TopicRegionValue};
358 use common_meta::peer::Peer;
359 use common_meta::rpc::router::{Region, RegionRoute};
360 use common_meta::wal_provider::RegionWalOptions;
361 use common_wal::options::KafkaWalOptions;
362 use store_api::storage::RegionId;
363
364 use super::*;
365 use crate::error::Error;
366 use crate::procedure::region_migration::test_util::{TestingEnv, new_procedure_context};
367 use crate::procedure::region_migration::{
368 ContextFactory, PersistentContext, RegionMigrationTriggerReason,
369 };
370 use crate::procedure::test_util::{
371 new_close_region_reply, new_upgrade_region_reply, send_mock_reply,
372 };
373
374 fn new_persistent_context() -> PersistentContext {
375 PersistentContext::new(
376 vec![("greptime".into(), "public".into())],
377 Peer::empty(1),
378 Peer::empty(2),
379 vec![RegionId::new(1024, 1)],
380 Duration::from_millis(1000),
381 RegionMigrationTriggerReason::Unknown,
382 )
383 }
384
385 fn kafka_wal_options(topic: &str) -> RegionWalOptions {
386 RegionWalOptions::from([(
387 1,
388 WalOptions::Kafka(KafkaWalOptions::new(topic.to_string())),
389 )])
390 }
391
392 async fn prepare_table_metadata(ctx: &Context, wal_options: RegionWalOptions) {
393 prepare_table_metadata_with_engine(ctx, wal_options, "engine").await;
394 }
395
396 async fn prepare_table_metadata_with_engine(
397 ctx: &Context,
398 wal_options: RegionWalOptions,
399 engine: &str,
400 ) {
401 let region_id = ctx.persistent_ctx.region_ids[0];
402 let mut table_info = new_test_table_info(region_id.table_id());
403 table_info.meta.engine = engine.to_string();
404 let region_routes = vec![RegionRoute {
405 region: Region::new_test(region_id),
406 leader_peer: Some(ctx.persistent_ctx.from_peer.clone()),
407 follower_peers: vec![ctx.persistent_ctx.to_peer.clone()],
408 ..Default::default()
409 }];
410 ctx.table_metadata_manager
411 .create_table_metadata(
412 table_info,
413 TableRouteValue::physical(region_routes),
414 wal_options,
415 )
416 .await
417 .unwrap();
418 }
419
420 #[tokio::test]
421 async fn test_build_upgrade_region_instruction_merges_topic_pruned_entry_id() {
422 let state = UpgradeCandidateRegion::default();
423 let persistent_context = new_persistent_context();
424 let env = TestingEnv::new();
425 let mut ctx = env.context_factory().new_context(persistent_context);
426 let region_id = ctx.persistent_ctx.region_ids[0];
427 let topic = "test_topic";
428 prepare_table_metadata(&ctx, kafka_wal_options(topic)).await;
429 ctx.table_metadata_manager
430 .topic_region_manager()
431 .batch_put(&[(
432 TopicRegionKey::new(region_id, topic),
433 Some(TopicRegionValue::new(Some(ReplayCheckpoint::new(10, None)))),
434 )])
435 .await
436 .unwrap();
437 ctx.table_metadata_manager
438 .topic_name_manager()
439 .batch_put(vec![TopicNameKey::new(topic)])
440 .await
441 .unwrap();
442 let prev = ctx
443 .table_metadata_manager
444 .topic_name_manager()
445 .get(topic)
446 .await
447 .unwrap();
448 ctx.table_metadata_manager
449 .topic_name_manager()
450 .update(topic, 20, prev)
451 .await
452 .unwrap();
453
454 let instruction = state
455 .build_upgrade_region_instruction(&mut ctx, Duration::from_secs(1))
456 .await
457 .unwrap();
458 let Instruction::UpgradeRegions(upgrade_regions) = instruction else {
459 unreachable!()
460 };
461
462 assert_eq!(upgrade_regions.len(), 1);
463 assert_eq!(upgrade_regions[0].replay_entry_id, Some(20));
464 assert_eq!(upgrade_regions[0].metadata_replay_entry_id, None);
465 }
466
467 #[tokio::test]
468 async fn test_build_upgrade_region_instruction_merges_metric_metadata_pruned_entry_id() {
469 let state = UpgradeCandidateRegion::default();
470 let persistent_context = new_persistent_context();
471 let env = TestingEnv::new();
472 let mut ctx = env.context_factory().new_context(persistent_context);
473 let region_id = ctx.persistent_ctx.region_ids[0];
474 let topic = "test_topic";
475 prepare_table_metadata_with_engine(&ctx, kafka_wal_options(topic), METRIC_ENGINE_NAME)
476 .await;
477 ctx.table_metadata_manager
478 .topic_region_manager()
479 .batch_put(&[(
480 TopicRegionKey::new(region_id, topic),
481 Some(TopicRegionValue::new(Some(ReplayCheckpoint::new(
482 10,
483 Some(5),
484 )))),
485 )])
486 .await
487 .unwrap();
488 ctx.table_metadata_manager
489 .topic_name_manager()
490 .batch_put(vec![TopicNameKey::new(topic)])
491 .await
492 .unwrap();
493 let prev = ctx
494 .table_metadata_manager
495 .topic_name_manager()
496 .get(topic)
497 .await
498 .unwrap();
499 ctx.table_metadata_manager
500 .topic_name_manager()
501 .update(topic, 20, prev)
502 .await
503 .unwrap();
504
505 let instruction = state
506 .build_upgrade_region_instruction(&mut ctx, Duration::from_secs(1))
507 .await
508 .unwrap();
509 let Instruction::UpgradeRegions(upgrade_regions) = instruction else {
510 unreachable!()
511 };
512
513 assert_eq!(upgrade_regions.len(), 1);
514 assert_eq!(upgrade_regions[0].replay_entry_id, Some(20));
515 assert_eq!(upgrade_regions[0].metadata_replay_entry_id, Some(20));
516 }
517
518 #[tokio::test]
519 async fn test_datanode_is_unreachable() {
520 let state = UpgradeCandidateRegion::default();
521 let persistent_context = new_persistent_context();
522 let env = TestingEnv::new();
523 let mut ctx = env.context_factory().new_context(persistent_context);
524 prepare_table_metadata(&ctx, HashMap::default()).await;
525 let err = state.upgrade_region(&mut ctx).await.unwrap_err();
526
527 assert_matches!(err, Error::PusherNotFound { .. });
528 assert!(!err.is_retryable());
529 }
530
531 #[tokio::test]
532 async fn test_pusher_dropped() {
533 let state = UpgradeCandidateRegion::default();
534 let persistent_context = new_persistent_context();
535 let to_peer_id = persistent_context.to_peer.id;
536
537 let mut env = TestingEnv::new();
538 let mut ctx = env.context_factory().new_context(persistent_context);
539 prepare_table_metadata(&ctx, HashMap::default()).await;
540 let mailbox_ctx = env.mailbox_context();
541
542 let (tx, rx) = tokio::sync::mpsc::channel(1);
543
544 mailbox_ctx
545 .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
546 .await;
547
548 drop(rx);
549
550 let err = state.upgrade_region(&mut ctx).await.unwrap_err();
551
552 assert_matches!(err, Error::PushMessage { .. });
553 assert!(!err.is_retryable());
554 }
555
556 #[tokio::test]
557 async fn test_procedure_exceeded_deadline() {
558 let state = UpgradeCandidateRegion::default();
559 let persistent_context = new_persistent_context();
560 let env = TestingEnv::new();
561 let mut ctx = env.context_factory().new_context(persistent_context);
562 prepare_table_metadata(&ctx, HashMap::default()).await;
563 ctx.volatile_ctx.metrics.operations_elapsed =
564 ctx.persistent_ctx.timeout + Duration::from_secs(1);
565
566 let err = state.upgrade_region(&mut ctx).await.unwrap_err();
567
568 assert_matches!(err, Error::ExceededDeadline { .. });
569 assert!(!err.is_retryable());
570 }
571
572 #[tokio::test]
573 async fn test_unexpected_instruction_reply() {
574 let state = UpgradeCandidateRegion::default();
575 let persistent_context = new_persistent_context();
576 let to_peer_id = persistent_context.to_peer.id;
577
578 let mut env = TestingEnv::new();
579 let mut ctx = env.context_factory().new_context(persistent_context);
580 prepare_table_metadata(&ctx, HashMap::default()).await;
581 let mailbox_ctx = env.mailbox_context();
582 let mailbox = mailbox_ctx.mailbox().clone();
583
584 let (tx, rx) = tokio::sync::mpsc::channel(1);
585
586 mailbox_ctx
587 .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
588 .await;
589
590 send_mock_reply(mailbox, rx, |id| Ok(new_close_region_reply(id)));
591
592 let err = state.upgrade_region(&mut ctx).await.unwrap_err();
593 assert_matches!(err, Error::UnexpectedInstructionReply { .. });
594 assert!(!err.is_retryable());
595 }
596
597 #[tokio::test]
598 async fn test_upgrade_region_failed() {
599 let state = UpgradeCandidateRegion::default();
600 let persistent_context = new_persistent_context();
601 let to_peer_id = persistent_context.to_peer.id;
602
603 let mut env = TestingEnv::new();
604 let mut ctx = env.context_factory().new_context(persistent_context);
605 prepare_table_metadata(&ctx, HashMap::default()).await;
606 let mailbox_ctx = env.mailbox_context();
607 let mailbox = mailbox_ctx.mailbox().clone();
608
609 let (tx, rx) = tokio::sync::mpsc::channel(1);
610
611 mailbox_ctx
612 .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
613 .await;
614
615 send_mock_reply(mailbox, rx, |id| {
617 Ok(new_upgrade_region_reply(
618 id,
619 true,
620 true,
621 Some("test mocked".to_string()),
622 ))
623 });
624
625 let err = state.upgrade_region(&mut ctx).await.unwrap_err();
626
627 assert_matches!(err, Error::RetryLater { .. });
628 assert!(err.is_retryable());
629 assert!(format!("{err:?}").contains("test mocked"));
630 }
631
632 #[tokio::test]
633 async fn test_upgrade_region_not_found() {
634 let state = UpgradeCandidateRegion::default();
635 let persistent_context = new_persistent_context();
636 let to_peer_id = persistent_context.to_peer.id;
637
638 let mut env = TestingEnv::new();
639 let mut ctx = env.context_factory().new_context(persistent_context);
640 prepare_table_metadata(&ctx, HashMap::default()).await;
641 let mailbox_ctx = env.mailbox_context();
642 let mailbox = mailbox_ctx.mailbox().clone();
643
644 let (tx, rx) = tokio::sync::mpsc::channel(1);
645
646 mailbox_ctx
647 .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
648 .await;
649
650 send_mock_reply(mailbox, rx, |id| {
651 Ok(new_upgrade_region_reply(id, true, false, None))
652 });
653
654 let err = state.upgrade_region(&mut ctx).await.unwrap_err();
655
656 assert_matches!(err, Error::Unexpected { .. });
657 assert!(!err.is_retryable());
658 assert!(err.to_string().contains("doesn't exist"));
659 }
660
661 #[tokio::test]
662 async fn test_upgrade_region_require_ready() {
663 let mut state = UpgradeCandidateRegion {
664 require_ready: true,
665 ..Default::default()
666 };
667
668 let persistent_context = new_persistent_context();
669 let to_peer_id = persistent_context.to_peer.id;
670
671 let mut env = TestingEnv::new();
672 let mut ctx = env.context_factory().new_context(persistent_context);
673 prepare_table_metadata(&ctx, HashMap::default()).await;
674 let mailbox_ctx = env.mailbox_context();
675 let mailbox = mailbox_ctx.mailbox().clone();
676
677 let (tx, rx) = tokio::sync::mpsc::channel(1);
678
679 mailbox_ctx
680 .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
681 .await;
682
683 send_mock_reply(mailbox, rx, |id| {
684 Ok(new_upgrade_region_reply(id, false, true, None))
685 });
686
687 let err = state.upgrade_region(&mut ctx).await.unwrap_err();
688
689 assert_matches!(err, Error::RetryLater { .. });
690 assert!(err.is_retryable());
691 assert!(format!("{err:?}").contains("still replaying the wal"));
692
693 state.require_ready = false;
695
696 let mailbox = mailbox_ctx.mailbox().clone();
697 let (tx, rx) = tokio::sync::mpsc::channel(1);
698
699 mailbox_ctx
700 .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
701 .await;
702
703 send_mock_reply(mailbox, rx, |id| {
704 Ok(new_upgrade_region_reply(id, false, true, None))
705 });
706
707 state.upgrade_region(&mut ctx).await.unwrap();
708 }
709
710 #[tokio::test]
711 async fn test_upgrade_region_with_retry_ok() {
712 let mut state = Box::<UpgradeCandidateRegion>::default();
713 state.retry_initial_interval = Duration::from_millis(100);
714 let persistent_context = new_persistent_context();
715 let to_peer_id = persistent_context.to_peer.id;
716
717 let mut env = TestingEnv::new();
718 let mut ctx = env.context_factory().new_context(persistent_context);
719 prepare_table_metadata(&ctx, HashMap::default()).await;
720 let mailbox_ctx = env.mailbox_context();
721 let mailbox = mailbox_ctx.mailbox().clone();
722
723 let (tx, mut rx) = tokio::sync::mpsc::channel(1);
724
725 mailbox_ctx
726 .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
727 .await;
728
729 common_runtime::spawn_global(async move {
730 let resp = rx.recv().await.unwrap().unwrap();
731 let reply_id = resp.mailbox_message.unwrap().id;
732 mailbox
733 .on_recv(
734 reply_id,
735 Err(error::MailboxTimeoutSnafu { id: reply_id }.build()),
736 )
737 .await
738 .unwrap();
739
740 let resp = rx.recv().await.unwrap().unwrap();
742 let reply_id = resp.mailbox_message.unwrap().id;
743 mailbox
744 .on_recv(
745 reply_id,
746 Ok(new_upgrade_region_reply(reply_id, false, true, None)),
747 )
748 .await
749 .unwrap();
750
751 let resp = rx.recv().await.unwrap().unwrap();
753 let reply_id = resp.mailbox_message.unwrap().id;
754 mailbox
755 .on_recv(
756 reply_id,
757 Ok(new_upgrade_region_reply(reply_id, true, true, None)),
758 )
759 .await
760 .unwrap();
761 });
762
763 let procedure_ctx = new_procedure_context();
764 let (next, _) = state.next(&mut ctx, &procedure_ctx).await.unwrap();
765
766 let update_metadata = next.as_any().downcast_ref::<UpdateMetadata>().unwrap();
767
768 assert_matches!(update_metadata, UpdateMetadata::Upgrade);
769 }
770
771 #[tokio::test]
772 async fn test_upgrade_region_with_retry_failed() {
773 let mut state = Box::<UpgradeCandidateRegion>::default();
774 state.retry_initial_interval = Duration::from_millis(100);
775 let persistent_context = new_persistent_context();
776 let to_peer_id = persistent_context.to_peer.id;
777
778 let mut env = TestingEnv::new();
779 let mut ctx = env.context_factory().new_context(persistent_context);
780 prepare_table_metadata(&ctx, HashMap::default()).await;
781 let mailbox_ctx = env.mailbox_context();
782 let mailbox = mailbox_ctx.mailbox().clone();
783
784 let (tx, mut rx) = tokio::sync::mpsc::channel(1);
785
786 mailbox_ctx
787 .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
788 .await;
789
790 common_runtime::spawn_global(async move {
791 let resp = rx.recv().await.unwrap().unwrap();
792 let reply_id = resp.mailbox_message.unwrap().id;
793 mailbox
794 .on_recv(
795 reply_id,
796 Err(error::MailboxTimeoutSnafu { id: reply_id }.build()),
797 )
798 .await
799 .unwrap();
800
801 let resp = rx.recv().await.unwrap().unwrap();
803 let reply_id = resp.mailbox_message.unwrap().id;
804 mailbox
805 .on_recv(
806 reply_id,
807 Ok(new_upgrade_region_reply(reply_id, false, true, None)),
808 )
809 .await
810 .unwrap();
811
812 let resp = rx.recv().await.unwrap().unwrap();
814 let reply_id = resp.mailbox_message.unwrap().id;
815 mailbox
816 .on_recv(
817 reply_id,
818 Ok(new_upgrade_region_reply(reply_id, false, false, None)),
819 )
820 .await
821 .unwrap();
822 });
823 let procedure_ctx = new_procedure_context();
824 let (next, _) = state.next(&mut ctx, &procedure_ctx).await.unwrap();
825
826 let update_metadata = next.as_any().downcast_ref::<UpdateMetadata>().unwrap();
827 assert_matches!(update_metadata, UpdateMetadata::Rollback);
828 }
829
830 #[tokio::test]
831 async fn test_upgrade_region_procedure_exceeded_deadline() {
832 let mut state = Box::<UpgradeCandidateRegion>::default();
833 state.retry_initial_interval = Duration::from_millis(100);
834 let persistent_context = new_persistent_context();
835 let to_peer_id = persistent_context.to_peer.id;
836
837 let mut env = TestingEnv::new();
838 let mut ctx = env.context_factory().new_context(persistent_context);
839 prepare_table_metadata(&ctx, HashMap::default()).await;
840 let mailbox_ctx = env.mailbox_context();
841 let mailbox = mailbox_ctx.mailbox().clone();
842 ctx.volatile_ctx.metrics.operations_elapsed =
843 ctx.persistent_ctx.timeout + Duration::from_secs(1);
844
845 let (tx, rx) = tokio::sync::mpsc::channel(1);
846 mailbox_ctx
847 .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
848 .await;
849
850 send_mock_reply(mailbox, rx, |id| {
851 Ok(new_upgrade_region_reply(id, false, true, None))
852 });
853 let procedure_ctx = new_procedure_context();
854 let (next, _) = state.next(&mut ctx, &procedure_ctx).await.unwrap();
855 let update_metadata = next.as_any().downcast_ref::<UpdateMetadata>().unwrap();
856 assert_matches!(update_metadata, UpdateMetadata::Rollback);
857 }
858}