Skip to main content

meta_srv/procedure/region_migration/
upgrade_candidate_region.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    // The optimistic retry times.
44    pub(crate) optimistic_retry: usize,
45    // The retry initial interval.
46    pub(crate) retry_initial_interval: Duration,
47    // If it's true it requires the candidate region MUST replay the WAL to the latest entry id.
48    // Otherwise, it will rollback to the old leader region.
49    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(&region_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    /// Builds upgrade region instruction.
118    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(&region_topic)
142            .await?;
143        // Build upgrade regions instruction.
144        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(&region_id)
150                .copied();
151            let metadata_last_entry_id = ctx
152                .volatile_ctx
153                .leader_region_metadata_last_entry_ids
154                .get(&region_id)
155                .copied();
156            let checkpoint = replay_checkpoints.get(&region_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    /// Tries to upgrade a candidate region.
222    ///
223    /// Retry:
224    /// - If `require_ready` is true, but the candidate region returns `ready` is false.
225    /// - [MailboxTimeout](error::Error::MailboxTimeout), Timeout.
226    ///
227    /// Abort:
228    /// - The candidate region doesn't exist.
229    /// - [PusherNotFound](error::Error::PusherNotFound), The datanode is unreachable.
230    /// - [PushMessage](error::Error::PushMessage), The receiver is dropped.
231    /// - [MailboxReceiver](error::Error::MailboxReceiver), The sender is dropped without sending (impossible).
232    /// - [UnexpectedInstructionReply](error::Error::UnexpectedInstructionReply) (impossible).
233    /// - [ExceededDeadline](error::Error::ExceededDeadline)
234    /// - Invalid JSON (impossible).
235    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    /// Upgrades a candidate region.
301    ///
302    /// Returns true if the candidate region is upgraded successfully.
303    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            // If using Kafka WAL, acquire a read lock on the topic to prevent WAL pruning during the upgrade.
316            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        // A reply contains an error.
616        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        // Sets the `require_ready` to false.
694        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            // retry: 1
741            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            // retry: 2
752            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            // retry: 1
802            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            // retry: 2
813            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}