datanode/heartbeat/handler/
sync_region.rs1use common_meta::instruction::{
16 InstructionError, InstructionReply, SyncRegion, SyncRegionReply, SyncRegionsReply,
17};
18use common_telemetry::{debug, error, info, warn};
19use futures::future::join_all;
20
21use crate::heartbeat::handler::{HandlerContext, InstructionHandler};
22
23#[derive(Debug, Clone, Copy, Default)]
26pub struct SyncRegionHandler;
27
28#[async_trait::async_trait]
29impl InstructionHandler for SyncRegionHandler {
30 type Instruction = Vec<SyncRegion>;
31
32 async fn handle(
34 &self,
35 ctx: &HandlerContext,
36 regions: Self::Instruction,
37 ) -> Option<InstructionReply> {
38 info!("Received sync region instructions: {:?}", regions);
39 let futures = regions
40 .into_iter()
41 .map(|sync_region| Self::handle_sync_region(ctx, sync_region));
42 let results = join_all(futures).await;
43 debug!("Sync region results: {:?}", results);
44
45 Some(InstructionReply::SyncRegions(SyncRegionsReply::new(
46 results,
47 )))
48 }
49}
50
51impl SyncRegionHandler {
52 async fn handle_sync_region(
54 ctx: &HandlerContext,
55 SyncRegion { region_id, request }: SyncRegion,
56 ) -> SyncRegionReply {
57 let Some(writable) = ctx.region_server.is_region_leader(region_id) else {
58 warn!("Region: {} is not found", region_id);
59 return SyncRegionReply {
60 region_id,
61 ready: false,
62 exists: false,
63 error: None,
64 };
65 };
66
67 if !writable {
68 warn!("Region: {} is not writable", region_id);
69 return SyncRegionReply {
70 region_id,
71 ready: false,
72 exists: true,
73 error: Some(InstructionError::legacy_internal_retryable(
74 "Region is not writable",
75 )),
76 };
77 }
78
79 match ctx.region_server.sync_region(region_id, request).await {
80 Ok(_) => {
81 info!("Successfully synced region: {}", region_id);
82 SyncRegionReply {
83 region_id,
84 ready: true,
85 exists: true,
86 error: None,
87 }
88 }
89 Err(e) => {
90 error!(e; "Failed to sync region: {}", region_id);
91 SyncRegionReply {
92 region_id,
93 ready: false,
94 exists: true,
95 error: Some(InstructionError::from_error(&e)),
96 }
97 }
98 }
99 }
100}
101
102#[cfg(test)]
103mod tests {
104 use std::sync::Arc;
105
106 use common_error::ext::{BoxedError, RetryHint};
107 use common_error::status_code::StatusCode;
108 use common_meta::kv_backend::memory::MemoryKvBackend;
109 use mito2::engine::MITO_ENGINE_NAME;
110 use mito2::error::ManifestDeltaNotFoundSnafu;
111 use object_store::{Error as ObjectStoreError, ErrorKind};
112 use snafu::IntoError;
113 use store_api::metric_engine_consts::METRIC_ENGINE_NAME;
114 use store_api::region_engine::{RegionRole, SyncRegionFromRequest};
115 use store_api::storage::RegionId;
116
117 use crate::error::HandleRegionRequestSnafu;
118 use crate::heartbeat::handler::sync_region::SyncRegionHandler;
119 use crate::heartbeat::handler::{HandlerContext, InstructionHandler};
120 use crate::tests::{MockRegionEngine, mock_region_server};
121
122 #[tokio::test]
123 async fn test_handle_sync_region_not_found() {
124 let mut mock_region_server = mock_region_server();
125 let (mock_engine, _) = MockRegionEngine::new(METRIC_ENGINE_NAME);
126 mock_region_server.register_engine(mock_engine);
127
128 let kv_backend = Arc::new(MemoryKvBackend::new());
129 let handler_context = HandlerContext::new_for_test(mock_region_server, kv_backend);
130 let handler = SyncRegionHandler;
131
132 let region_id = RegionId::new(1024, 1);
133 let sync_region = common_meta::instruction::SyncRegion {
134 region_id,
135 request: SyncRegionFromRequest::from_manifest(Default::default()),
136 };
137
138 let reply = handler
139 .handle(&handler_context, vec![sync_region])
140 .await
141 .unwrap()
142 .expect_sync_regions_reply();
143
144 assert_eq!(reply.len(), 1);
145 assert_eq!(reply[0].region_id, region_id);
146 assert!(!reply[0].exists);
147 assert!(!reply[0].ready);
148 }
149
150 #[tokio::test]
151 async fn test_handle_sync_region_not_writable() {
152 let mock_region_server = mock_region_server();
153 let region_id = RegionId::new(1024, 1);
154 let (mock_engine, _) = MockRegionEngine::with_custom_apply_fn(METRIC_ENGINE_NAME, |r| {
155 r.mock_role = Some(Some(RegionRole::Follower));
156 });
157 mock_region_server.register_test_region(region_id, mock_engine);
158
159 let kv_backend = Arc::new(MemoryKvBackend::new());
160 let handler_context = HandlerContext::new_for_test(mock_region_server, kv_backend);
161 let handler = SyncRegionHandler;
162
163 let sync_region = common_meta::instruction::SyncRegion {
164 region_id,
165 request: SyncRegionFromRequest::from_manifest(Default::default()),
166 };
167
168 let reply = handler
169 .handle(&handler_context, vec![sync_region])
170 .await
171 .unwrap()
172 .expect_sync_regions_reply();
173
174 assert_eq!(reply.len(), 1);
175 assert_eq!(reply[0].region_id, region_id);
176 assert!(reply[0].exists);
177 assert!(!reply[0].ready);
178 assert!(reply[0].error.is_some());
179 }
180
181 #[tokio::test]
182 async fn test_handle_sync_region_success() {
183 let mock_region_server = mock_region_server();
184 let region_id = RegionId::new(1024, 1);
185 let (mock_engine, _) = MockRegionEngine::with_custom_apply_fn(METRIC_ENGINE_NAME, |r| {
186 r.mock_role = Some(Some(RegionRole::Leader));
187 });
188 mock_region_server.register_test_region(region_id, mock_engine);
189
190 let kv_backend = Arc::new(MemoryKvBackend::new());
191 let handler_context = HandlerContext::new_for_test(mock_region_server, kv_backend);
192 let handler = SyncRegionHandler;
193
194 let sync_region = common_meta::instruction::SyncRegion {
195 region_id,
196 request: SyncRegionFromRequest::from_manifest(Default::default()),
197 };
198
199 let reply = handler
200 .handle(&handler_context, vec![sync_region])
201 .await
202 .unwrap()
203 .expect_sync_regions_reply();
204
205 assert_eq!(reply.len(), 1);
206 assert_eq!(reply[0].region_id, region_id);
207 assert!(reply[0].exists);
208 assert!(reply[0].ready);
209 assert!(reply[0].error.is_none());
210 }
211
212 #[tokio::test]
213 async fn test_handle_sync_region_preserves_manifest_delta_not_found_retry_hint() {
214 let mock_region_server = mock_region_server();
215 let region_id = RegionId::new(1024, 1);
216 let (mock_engine, _) = MockRegionEngine::with_custom_apply_fn(MITO_ENGINE_NAME, |engine| {
217 engine.mock_role = Some(Some(RegionRole::Leader));
218 engine.handle_sync_region_mock_fn = Some(Box::new(|region_id, _request| {
219 let manifest_error = ManifestDeltaNotFoundSnafu {
220 version: 1_u64,
221 path: "manifest/00000000000000000001.json",
222 }
223 .into_error(ObjectStoreError::new(
224 ErrorKind::NotFound,
225 "mock listed manifest delta not found",
226 ));
227 Err(HandleRegionRequestSnafu { region_id }
228 .into_error(BoxedError::new(manifest_error)))
229 }));
230 });
231 mock_region_server.register_test_region(region_id, mock_engine);
232
233 let handler_context =
234 HandlerContext::new_for_test(mock_region_server, Arc::new(MemoryKvBackend::new()));
235 let sync_region = common_meta::instruction::SyncRegion {
236 region_id,
237 request: SyncRegionFromRequest::from_manifest(Default::default()),
238 };
239
240 let reply = SyncRegionHandler
241 .handle(&handler_context, vec![sync_region])
242 .await
243 .unwrap()
244 .expect_sync_regions_reply();
245
246 assert_eq!(1, reply.len());
247 assert!(reply[0].exists);
248 assert!(!reply[0].ready);
249 let error = reply[0].error.as_ref().unwrap();
250 assert_eq!(StatusCode::StorageUnavailable, error.code);
251 assert_eq!(RetryHint::Retryable, error.retry_hint);
252 assert!(error.message.contains("00000000000000000001.json"));
253 }
254}