Skip to main content

meta_srv/procedure/region_migration/
downgrade_leader_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::time::Duration;
17
18use api::v1::meta::MailboxMessage;
19use common_error::ext::BoxedError;
20use common_meta::distributed_time_constants::default_distributed_time_constants;
21use common_meta::instruction::{
22    DowngradeRegion, DowngradeRegionReply, DowngradeRegionsReply, Instruction, InstructionReply,
23};
24use common_procedure::{Context as ProcedureContext, Status};
25use common_telemetry::tracing_context::TracingContext;
26use common_telemetry::{debug, error, info, warn};
27use common_time::util::current_time_millis;
28use serde::{Deserialize, Serialize};
29use snafu::{OptionExt, ResultExt};
30use tokio::time::{Instant, sleep};
31
32use crate::discovery::utils::find_datanode_lease_value;
33use crate::error::{self, Result};
34use crate::handler::HeartbeatMailbox;
35use crate::procedure::region_migration::update_metadata::UpdateMetadata;
36use crate::procedure::region_migration::upgrade_candidate_region::UpgradeCandidateRegion;
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 DowngradeLeaderRegion {
43    // The optimistic retry times.
44    optimistic_retry: usize,
45    // The retry initial interval.
46    retry_initial_interval: Duration,
47}
48
49impl Default for DowngradeLeaderRegion {
50    fn default() -> Self {
51        Self {
52            optimistic_retry: 3,
53            retry_initial_interval: Duration::from_millis(500),
54        }
55    }
56}
57
58#[async_trait::async_trait]
59#[typetag::serde]
60impl State for DowngradeLeaderRegion {
61    async fn next(
62        &mut self,
63        ctx: &mut Context,
64        _procedure_ctx: &ProcedureContext,
65    ) -> Result<(Box<dyn State>, Status)> {
66        let now = Instant::now();
67        // Ensures the `leader_region_lease_deadline` must exist after recovering.
68        ctx.volatile_ctx
69            .set_leader_region_lease_deadline(default_distributed_time_constants().region_lease);
70
71        match self.downgrade_region_with_retry(ctx).await {
72            Ok(_) => {
73                // Do nothing
74                info!(
75                    "Downgraded region leader success, region: {:?}",
76                    ctx.persistent_ctx.region_ids
77                );
78            }
79            Err(error::Error::ExceededDeadline { .. }) => {
80                info!(
81                    "Downgrade region leader exceeded deadline, region: {:?}",
82                    ctx.persistent_ctx.region_ids
83                );
84                // Rollbacks the metadata if procedure is timeout
85                return Ok((Box::new(UpdateMetadata::Rollback), Status::executing(false)));
86            }
87            Err(err) => {
88                error!(err; "Occurs non-retryable error, region: {:?}", ctx.persistent_ctx.region_ids);
89                if let Some(deadline) = ctx.volatile_ctx.leader_region_lease_deadline.as_ref() {
90                    info!(
91                        "Running into the downgrade region leader slow path, region: {:?}, sleep until {:?}",
92                        ctx.persistent_ctx.region_ids, deadline
93                    );
94                    tokio::time::sleep_until(*deadline).await;
95                } else {
96                    warn!(
97                        "Leader region lease deadline is not set, region: {:?}",
98                        ctx.persistent_ctx.region_ids
99                    );
100                }
101            }
102        }
103        ctx.update_downgrade_leader_region_elapsed(now);
104
105        Ok((
106            Box::new(UpgradeCandidateRegion::default()),
107            Status::executing(false),
108        ))
109    }
110
111    fn as_any(&self) -> &dyn Any {
112        self
113    }
114}
115
116impl DowngradeLeaderRegion {
117    /// Builds downgrade region instruction.
118    fn build_downgrade_region_instruction(
119        &self,
120        ctx: &Context,
121        flush_timeout: Duration,
122    ) -> Instruction {
123        let region_ids = &ctx.persistent_ctx.region_ids;
124        let mut downgrade_regions = Vec::with_capacity(region_ids.len());
125        for region_id in region_ids {
126            downgrade_regions.push(DowngradeRegion {
127                region_id: *region_id,
128                flush_timeout: Some(flush_timeout),
129            });
130        }
131
132        Instruction::DowngradeRegions(downgrade_regions)
133    }
134
135    fn handle_downgrade_region_reply(
136        &self,
137        ctx: &mut Context,
138        reply: &DowngradeRegionReply,
139        now: &Instant,
140    ) -> Result<()> {
141        let leader = &ctx.persistent_ctx.from_peer;
142        let DowngradeRegionReply {
143            region_id,
144            last_entry_id,
145            metadata_last_entry_id,
146            exists,
147            error,
148        } = reply;
149
150        if let Some(error) = error {
151            return instruction_error_result(
152                error,
153                format!(
154                    "Failed to downgrade the region {} on datanode {:?}, error: {:?}, elapsed: {:?}",
155                    region_id,
156                    leader,
157                    error,
158                    now.elapsed()
159                ),
160            );
161        }
162
163        if !exists {
164            warn!(
165                "Trying to downgrade the region {} on datanode {:?}, but region doesn't exist!, elapsed: {:?}",
166                region_id,
167                leader,
168                now.elapsed()
169            );
170        } else {
171            info!(
172                "Region {} leader is downgraded on datanode {:?}, last_entry_id: {:?}, metadata_last_entry_id: {:?}, elapsed: {:?}",
173                region_id,
174                leader,
175                last_entry_id,
176                metadata_last_entry_id,
177                now.elapsed()
178            );
179        }
180
181        if let Some(last_entry_id) = last_entry_id {
182            debug!(
183                "set last_entry_id: {:?}, region_id: {:?}",
184                last_entry_id, region_id
185            );
186            ctx.volatile_ctx
187                .set_last_entry_id(*region_id, *last_entry_id);
188        }
189
190        if let Some(metadata_last_entry_id) = metadata_last_entry_id {
191            ctx.volatile_ctx
192                .set_metadata_last_entry_id(*region_id, *metadata_last_entry_id);
193        }
194
195        Ok(())
196    }
197
198    /// Tries to downgrade a leader region.
199    ///
200    /// Retry:
201    /// - [MailboxTimeout](error::Error::MailboxTimeout), Timeout.
202    /// - Failed to downgrade region on the Datanode.
203    ///
204    /// Abort:
205    /// - [PusherNotFound](error::Error::PusherNotFound), The datanode is unreachable.
206    /// - [PushMessage](error::Error::PushMessage), The receiver is dropped.
207    /// - [MailboxReceiver](error::Error::MailboxReceiver), The sender is dropped without sending (impossible).
208    /// - [UnexpectedInstructionReply](error::Error::UnexpectedInstructionReply).
209    /// - [ExceededDeadline](error::Error::ExceededDeadline)
210    /// - Invalid JSON.
211    async fn downgrade_region(&self, ctx: &mut Context) -> Result<()> {
212        let region_ids = &ctx.persistent_ctx.region_ids;
213        let operation_timeout =
214            ctx.next_operation_timeout()
215                .context(error::ExceededDeadlineSnafu {
216                    operation: "Downgrade region",
217                })?;
218        let downgrade_instruction = self.build_downgrade_region_instruction(ctx, operation_timeout);
219
220        let leader = &ctx.persistent_ctx.from_peer;
221        let tracing_ctx = TracingContext::from_current_span();
222        let msg = MailboxMessage::json_message(
223            &format!("Downgrade leader regions: {:?}", region_ids),
224            &format!("Metasrv@{}", ctx.server_addr()),
225            &format!("Datanode-{}@{}", leader.id, leader.addr),
226            common_time::util::current_time_millis(),
227            &downgrade_instruction,
228            Some(tracing_ctx.to_w3c()),
229        )
230        .with_context(|_| error::SerializeToJsonSnafu {
231            input: downgrade_instruction.to_string(),
232        })?;
233
234        let ch = Channel::Datanode(leader.id);
235        let now = Instant::now();
236        let receiver = ctx.mailbox.send(&ch, msg, operation_timeout).await?;
237
238        match receiver.await {
239            Ok(msg) => {
240                let reply = HeartbeatMailbox::json_reply(&msg)?;
241                info!(
242                    "Received downgrade region reply: {:?}, region: {:?}, elapsed: {:?}",
243                    reply,
244                    region_ids,
245                    now.elapsed()
246                );
247                let InstructionReply::DowngradeRegions(DowngradeRegionsReply { replies }) = reply
248                else {
249                    return error::UnexpectedInstructionReplySnafu {
250                        mailbox_message: msg.to_string(),
251                        reason: "expect downgrade region reply",
252                    }
253                    .fail();
254                };
255
256                for reply in replies {
257                    self.handle_downgrade_region_reply(ctx, &reply, &now)?;
258                }
259                Ok(())
260            }
261            Err(error::Error::MailboxTimeout { .. }) => {
262                let reason = format!(
263                    "Mailbox received timeout for downgrade leader region {region_ids:?} on datanode {:?}, elapsed: {:?}",
264                    leader,
265                    now.elapsed()
266                );
267                error::RetryLaterSnafu { reason }.fail()
268            }
269            Err(err) => Err(err),
270        }
271    }
272
273    async fn update_leader_region_lease_deadline(&self, ctx: &mut Context) {
274        let leader = &ctx.persistent_ctx.from_peer;
275
276        let last_connection_at = match find_datanode_lease_value(&ctx.in_memory, leader.id).await {
277            Ok(lease_value) => lease_value.map(|lease_value| lease_value.timestamp_millis),
278            Err(err) => {
279                error!(err; "Failed to find datanode lease value for datanode: {}, during region migration, region: {:?}", leader, ctx.persistent_ctx.region_ids);
280                return;
281            }
282        };
283
284        if let Some(last_connection_at) = last_connection_at {
285            let now = current_time_millis();
286            let elapsed = now - last_connection_at;
287            let region_lease = default_distributed_time_constants().region_lease;
288
289            // It's safe to update the region leader lease deadline here because:
290            // 1. The old region leader has already been marked as downgraded in metadata,
291            //    which means any attempts to renew its lease will be rejected.
292            // 2. The pusher disconnect time record only gets removed when the datanode (from_peer)
293            //    establishes a new heartbeat connection stream.
294            if elapsed >= (region_lease.as_secs() * 1000) as i64 {
295                ctx.volatile_ctx.reset_leader_region_lease_deadline();
296                info!(
297                    "Datanode {}({}) has been disconnected for longer than the region lease period ({:?}), reset leader region lease deadline to None, region: {:?}",
298                    leader, last_connection_at, region_lease, ctx.persistent_ctx.region_ids
299                );
300            } else if elapsed > 0 {
301                // `now - last_connection_at` < REGION_LEASE_SECS * 1000
302                let lease_timeout =
303                    region_lease - Duration::from_millis((now - last_connection_at) as u64);
304                ctx.volatile_ctx.reset_leader_region_lease_deadline();
305                ctx.volatile_ctx
306                    .set_leader_region_lease_deadline(lease_timeout);
307                info!(
308                    "Datanode {}({}) last connected {:?} ago, updated leader region lease deadline to {:?}, region: {:?}",
309                    leader,
310                    last_connection_at,
311                    elapsed,
312                    ctx.volatile_ctx.leader_region_lease_deadline,
313                    ctx.persistent_ctx.region_ids
314                );
315            } else {
316                warn!(
317                    "Datanode {} has invalid last connection timestamp: {} (which is after current time: {}), region: {:?}",
318                    leader, last_connection_at, now, ctx.persistent_ctx.region_ids
319                )
320            }
321        } else {
322            warn!(
323                "Failed to find last connection time for datanode {}, unable to update region lease deadline, region: {:?}",
324                leader, ctx.persistent_ctx.region_ids
325            )
326        }
327    }
328
329    /// Downgrades a leader region.
330    ///
331    /// Fast path:
332    /// - Waits for the reply of downgrade instruction.
333    ///
334    /// Slow path:
335    /// - Waits for the lease of the leader region expired.
336    ///
337    /// Abort:
338    /// - ExceededDeadline
339    async fn downgrade_region_with_retry(&self, ctx: &mut Context) -> Result<()> {
340        let mut retry = 0;
341
342        loop {
343            let timer = Instant::now();
344            if let Err(err) = self.downgrade_region(ctx).await {
345                ctx.update_operations_elapsed(timer);
346                retry += 1;
347                // Throws the error immediately if the procedure exceeded the deadline.
348                if matches!(err, error::Error::ExceededDeadline { .. }) {
349                    error!(err; "Failed to downgrade region leader, regions: {:?}, exceeded deadline", ctx.persistent_ctx.region_ids);
350                    return Err(err);
351                } else if matches!(err, error::Error::PusherNotFound { .. }) {
352                    // Throws the error immediately if the datanode is unreachable.
353                    error!(err; "Failed to downgrade region leader, regions: {:?}, datanode({}) is unreachable(PusherNotFound)", ctx.persistent_ctx.region_ids, ctx.persistent_ctx.from_peer.id);
354                    self.update_leader_region_lease_deadline(ctx).await;
355                    return Err(err);
356                } else if err.is_retryable() && retry < self.optimistic_retry {
357                    error!(err; "Failed to downgrade region leader, regions: {:?}, retry later", ctx.persistent_ctx.region_ids);
358                    sleep(self.retry_initial_interval).await;
359                } else {
360                    return Err(BoxedError::new(err)).context(error::DowngradeLeaderSnafu {
361                        // TODO(weny): handle multiple regions.
362                        region_id: ctx.persistent_ctx.region_ids[0],
363                    })?;
364                }
365            } else {
366                ctx.update_operations_elapsed(timer);
367                // Resets the deadline.
368                ctx.volatile_ctx.reset_leader_region_lease_deadline();
369                break;
370            }
371        }
372
373        Ok(())
374    }
375}
376
377#[cfg(test)]
378mod tests {
379    use std::assert_matches;
380    use std::collections::HashMap;
381
382    use common_meta::key::table_route::TableRouteValue;
383    use common_meta::key::test_utils::new_test_table_info;
384    use common_meta::peer::Peer;
385    use common_meta::rpc::router::{Region, RegionRoute};
386    use common_meta::wal_provider::RegionWalOptions;
387    use store_api::storage::RegionId;
388    use tokio::time::Instant;
389
390    use super::*;
391    use crate::error::Error;
392    use crate::procedure::region_migration::test_util::{TestingEnv, new_procedure_context};
393    use crate::procedure::region_migration::{
394        ContextFactory, PersistentContext, RegionMigrationTriggerReason,
395    };
396    use crate::procedure::test_util::{
397        new_close_region_reply, new_downgrade_region_reply, send_mock_reply,
398    };
399
400    fn new_persistent_context() -> PersistentContext {
401        PersistentContext::new(
402            vec![("greptime".into(), "public".into())],
403            Peer::empty(1),
404            Peer::empty(2),
405            vec![RegionId::new(1024, 1)],
406            Duration::from_millis(1000),
407            RegionMigrationTriggerReason::Unknown,
408        )
409    }
410
411    async fn prepare_table_metadata(ctx: &Context, wal_options: RegionWalOptions) {
412        let region_id = ctx.persistent_ctx.region_ids[0];
413        let table_info = new_test_table_info(region_id.table_id());
414        let region_routes = vec![RegionRoute {
415            region: Region::new_test(region_id),
416            leader_peer: Some(ctx.persistent_ctx.from_peer.clone()),
417            follower_peers: vec![ctx.persistent_ctx.to_peer.clone()],
418            ..Default::default()
419        }];
420        ctx.table_metadata_manager
421            .create_table_metadata(
422                table_info,
423                TableRouteValue::physical(region_routes),
424                wal_options,
425            )
426            .await
427            .unwrap();
428    }
429
430    #[tokio::test]
431    async fn test_datanode_is_unreachable() {
432        let state = DowngradeLeaderRegion::default();
433        let persistent_context = new_persistent_context();
434        let env = TestingEnv::new();
435        let mut ctx = env.context_factory().new_context(persistent_context);
436        prepare_table_metadata(&ctx, HashMap::default()).await;
437        let err = state.downgrade_region(&mut ctx).await.unwrap_err();
438
439        assert_matches!(err, Error::PusherNotFound { .. });
440        assert!(!err.is_retryable());
441    }
442
443    #[tokio::test]
444    async fn test_pusher_dropped() {
445        let state = DowngradeLeaderRegion::default();
446        let persistent_context = new_persistent_context();
447        let from_peer_id = persistent_context.from_peer.id;
448
449        let mut env = TestingEnv::new();
450        let mut ctx = env.context_factory().new_context(persistent_context);
451        prepare_table_metadata(&ctx, HashMap::default()).await;
452        let mailbox_ctx = env.mailbox_context();
453
454        let (tx, rx) = tokio::sync::mpsc::channel(1);
455
456        mailbox_ctx
457            .insert_heartbeat_response_receiver(Channel::Datanode(from_peer_id), tx)
458            .await;
459
460        drop(rx);
461
462        let err = state.downgrade_region(&mut ctx).await.unwrap_err();
463
464        assert_matches!(err, Error::PushMessage { .. });
465        assert!(!err.is_retryable());
466    }
467
468    #[tokio::test]
469    async fn test_procedure_exceeded_deadline() {
470        let state = DowngradeLeaderRegion::default();
471        let persistent_context = new_persistent_context();
472        let env = TestingEnv::new();
473        let mut ctx = env.context_factory().new_context(persistent_context);
474        prepare_table_metadata(&ctx, HashMap::default()).await;
475        ctx.volatile_ctx.metrics.operations_elapsed =
476            ctx.persistent_ctx.timeout + Duration::from_secs(1);
477
478        let err = state.downgrade_region(&mut ctx).await.unwrap_err();
479
480        assert_matches!(err, Error::ExceededDeadline { .. });
481        assert!(!err.is_retryable());
482
483        let err = state
484            .downgrade_region_with_retry(&mut ctx)
485            .await
486            .unwrap_err();
487        assert_matches!(err, Error::ExceededDeadline { .. });
488        assert!(!err.is_retryable());
489    }
490
491    #[tokio::test]
492    async fn test_unexpected_instruction_reply() {
493        let state = DowngradeLeaderRegion::default();
494        let persistent_context = new_persistent_context();
495        let from_peer_id = persistent_context.from_peer.id;
496
497        let mut env = TestingEnv::new();
498        let mut ctx = env.context_factory().new_context(persistent_context);
499        prepare_table_metadata(&ctx, HashMap::default()).await;
500        let mailbox_ctx = env.mailbox_context();
501        let mailbox = mailbox_ctx.mailbox().clone();
502
503        let (tx, rx) = tokio::sync::mpsc::channel(1);
504
505        mailbox_ctx
506            .insert_heartbeat_response_receiver(Channel::Datanode(from_peer_id), tx)
507            .await;
508
509        // Sends an incorrect reply.
510        send_mock_reply(mailbox, rx, |id| Ok(new_close_region_reply(id)));
511
512        let err = state.downgrade_region(&mut ctx).await.unwrap_err();
513
514        assert_matches!(err, Error::UnexpectedInstructionReply { .. });
515        assert!(!err.is_retryable());
516    }
517
518    #[tokio::test]
519    async fn test_instruction_exceeded_deadline() {
520        let state = DowngradeLeaderRegion::default();
521        let persistent_context = new_persistent_context();
522        let from_peer_id = persistent_context.from_peer.id;
523
524        let mut env = TestingEnv::new();
525        let mut ctx = env.context_factory().new_context(persistent_context);
526        prepare_table_metadata(&ctx, HashMap::default()).await;
527        let mailbox_ctx = env.mailbox_context();
528        let mailbox = mailbox_ctx.mailbox().clone();
529
530        let (tx, rx) = tokio::sync::mpsc::channel(1);
531
532        mailbox_ctx
533            .insert_heartbeat_response_receiver(Channel::Datanode(from_peer_id), tx)
534            .await;
535
536        send_mock_reply(mailbox, rx, |id| {
537            Err(error::MailboxTimeoutSnafu { id }.build())
538        });
539
540        let err = state.downgrade_region(&mut ctx).await.unwrap_err();
541
542        assert_matches!(err, Error::RetryLater { .. });
543        assert!(err.is_retryable());
544    }
545
546    #[tokio::test]
547    async fn test_downgrade_region_failed() {
548        let state = DowngradeLeaderRegion::default();
549        let persistent_context = new_persistent_context();
550        let from_peer_id = persistent_context.from_peer.id;
551
552        let mut env = TestingEnv::new();
553        let mut ctx = env.context_factory().new_context(persistent_context);
554        prepare_table_metadata(&ctx, HashMap::default()).await;
555        let mailbox_ctx = env.mailbox_context();
556        let mailbox = mailbox_ctx.mailbox().clone();
557
558        let (tx, rx) = tokio::sync::mpsc::channel(1);
559
560        mailbox_ctx
561            .insert_heartbeat_response_receiver(Channel::Datanode(from_peer_id), tx)
562            .await;
563
564        send_mock_reply(mailbox, rx, |id| {
565            Ok(new_downgrade_region_reply(
566                id,
567                None,
568                false,
569                Some("test mocked".to_string()),
570            ))
571        });
572
573        let err = state.downgrade_region(&mut ctx).await.unwrap_err();
574
575        assert_matches!(err, Error::RetryLater { .. });
576        assert!(err.is_retryable());
577        assert!(format!("{err:?}").contains("test mocked"), "err: {err:?}",);
578    }
579
580    #[tokio::test]
581    async fn test_downgrade_region_with_retry_fast_path() {
582        let state = DowngradeLeaderRegion::default();
583        let persistent_context = new_persistent_context();
584        let from_peer_id = persistent_context.from_peer.id;
585
586        let mut env = TestingEnv::new();
587        let mut ctx = env.context_factory().new_context(persistent_context);
588        prepare_table_metadata(&ctx, HashMap::default()).await;
589        let mailbox_ctx = env.mailbox_context();
590        let mailbox = mailbox_ctx.mailbox().clone();
591
592        let (tx, mut rx) = tokio::sync::mpsc::channel(1);
593
594        mailbox_ctx
595            .insert_heartbeat_response_receiver(Channel::Datanode(from_peer_id), tx)
596            .await;
597
598        common_runtime::spawn_global(async move {
599            // retry: 0.
600            let resp = rx.recv().await.unwrap().unwrap();
601            let reply_id = resp.mailbox_message.unwrap().id;
602            mailbox
603                .on_recv(
604                    reply_id,
605                    Err(error::MailboxTimeoutSnafu { id: reply_id }.build()),
606                )
607                .await
608                .unwrap();
609
610            // retry: 1.
611            let resp = rx.recv().await.unwrap().unwrap();
612            let reply_id = resp.mailbox_message.unwrap().id;
613            mailbox
614                .on_recv(
615                    reply_id,
616                    Ok(new_downgrade_region_reply(reply_id, Some(1), true, None)),
617                )
618                .await
619                .unwrap();
620        });
621
622        state.downgrade_region_with_retry(&mut ctx).await.unwrap();
623        assert_eq!(
624            ctx.volatile_ctx
625                .leader_region_last_entry_ids
626                .get(&RegionId::new(0, 0))
627                .cloned(),
628            Some(1)
629        );
630        assert!(ctx.volatile_ctx.leader_region_lease_deadline.is_none());
631    }
632
633    #[tokio::test]
634    async fn test_downgrade_region_with_retry_slow_path() {
635        let state = DowngradeLeaderRegion {
636            optimistic_retry: 3,
637            retry_initial_interval: Duration::from_millis(100),
638        };
639        let persistent_context = new_persistent_context();
640        let from_peer_id = persistent_context.from_peer.id;
641
642        let mut env = TestingEnv::new();
643        let mut ctx = env.context_factory().new_context(persistent_context);
644        let mailbox_ctx = env.mailbox_context();
645        let mailbox = mailbox_ctx.mailbox().clone();
646
647        let (tx, mut rx) = tokio::sync::mpsc::channel(1);
648
649        mailbox_ctx
650            .insert_heartbeat_response_receiver(Channel::Datanode(from_peer_id), tx)
651            .await;
652
653        common_runtime::spawn_global(async move {
654            for _ in 0..3 {
655                let resp = rx.recv().await.unwrap().unwrap();
656                let reply_id = resp.mailbox_message.unwrap().id;
657                mailbox
658                    .on_recv(
659                        reply_id,
660                        Err(error::MailboxTimeoutSnafu { id: reply_id }.build()),
661                    )
662                    .await
663                    .unwrap();
664            }
665        });
666
667        ctx.volatile_ctx
668            .set_leader_region_lease_deadline(Duration::from_secs(5));
669        let expected_deadline = ctx.volatile_ctx.leader_region_lease_deadline.unwrap();
670        let err = state
671            .downgrade_region_with_retry(&mut ctx)
672            .await
673            .unwrap_err();
674        assert_matches!(err, error::Error::DowngradeLeader { .. });
675        // assert_eq!(ctx.volatile_ctx.leader_region_last_entry_id, None);
676        // Should remain no change.
677        assert_eq!(
678            ctx.volatile_ctx.leader_region_lease_deadline.unwrap(),
679            expected_deadline
680        )
681    }
682
683    #[tokio::test]
684    async fn test_next_upgrade_candidate_state() {
685        let mut state = Box::<DowngradeLeaderRegion>::default();
686        let persistent_context = new_persistent_context();
687        let from_peer_id = persistent_context.from_peer.id;
688
689        let mut env = TestingEnv::new();
690        let mut ctx = env.context_factory().new_context(persistent_context);
691        prepare_table_metadata(&ctx, HashMap::default()).await;
692        let mailbox_ctx = env.mailbox_context();
693        let mailbox = mailbox_ctx.mailbox().clone();
694
695        let (tx, rx) = tokio::sync::mpsc::channel(1);
696
697        mailbox_ctx
698            .insert_heartbeat_response_receiver(Channel::Datanode(from_peer_id), tx)
699            .await;
700
701        send_mock_reply(mailbox, rx, |id| {
702            Ok(new_downgrade_region_reply(id, Some(1), true, None))
703        });
704
705        let timer = Instant::now();
706        let procedure_ctx = new_procedure_context();
707        let (next, _) = state.next(&mut ctx, &procedure_ctx).await.unwrap();
708        let elapsed = timer.elapsed().as_secs();
709        let region_lease = default_distributed_time_constants().region_lease.as_secs();
710        assert!(elapsed < region_lease / 2);
711        assert_eq!(
712            ctx.volatile_ctx
713                .leader_region_last_entry_ids
714                .get(&RegionId::new(0, 0))
715                .cloned(),
716            Some(1)
717        );
718        assert!(ctx.volatile_ctx.leader_region_lease_deadline.is_none());
719
720        let _ = next
721            .as_any()
722            .downcast_ref::<UpgradeCandidateRegion>()
723            .unwrap();
724    }
725
726    #[tokio::test]
727    async fn test_downgrade_region_procedure_exceeded_deadline() {
728        let mut state = Box::<UpgradeCandidateRegion>::default();
729        state.retry_initial_interval = Duration::from_millis(100);
730        let persistent_context = new_persistent_context();
731        let to_peer_id = persistent_context.to_peer.id;
732
733        let mut env = TestingEnv::new();
734        let mut ctx = env.context_factory().new_context(persistent_context);
735        let mailbox_ctx = env.mailbox_context();
736        let mailbox = mailbox_ctx.mailbox().clone();
737        ctx.volatile_ctx.metrics.operations_elapsed =
738            ctx.persistent_ctx.timeout + Duration::from_secs(1);
739
740        let (tx, rx) = tokio::sync::mpsc::channel(1);
741        mailbox_ctx
742            .insert_heartbeat_response_receiver(Channel::Datanode(to_peer_id), tx)
743            .await;
744
745        send_mock_reply(mailbox, rx, |id| {
746            Ok(new_downgrade_region_reply(id, None, true, None))
747        });
748        let procedure_ctx = new_procedure_context();
749        let (next, _) = state.next(&mut ctx, &procedure_ctx).await.unwrap();
750        let update_metadata = next.as_any().downcast_ref::<UpdateMetadata>().unwrap();
751        assert_matches!(update_metadata, UpdateMetadata::Rollback);
752    }
753}