1use 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 optimistic_retry: usize,
45 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 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 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 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 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 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 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 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 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 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 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 region_id: ctx.persistent_ctx.region_ids[0],
363 })?;
364 }
365 } else {
366 ctx.update_operations_elapsed(timer);
367 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 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 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 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!(
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}