datanode/heartbeat/handler/
open_region.rs1use common_meta::instruction::{InstructionError, InstructionReply, OpenRegion, SimpleReply};
16use common_meta::wal_provider::serialize_wal_options;
17use common_telemetry::info;
18use store_api::path_utils::table_dir;
19use store_api::region_request::{PathType, RegionOpenRequest};
20use store_api::storage::RegionId;
21
22use crate::heartbeat::handler::{HandlerContext, InstructionHandler};
23
24pub struct OpenRegionsHandler {
25 pub open_region_parallelism: usize,
26}
27
28#[async_trait::async_trait]
29impl InstructionHandler for OpenRegionsHandler {
30 type Instruction = Vec<OpenRegion>;
31 async fn handle(
32 &self,
33 ctx: &HandlerContext,
34 open_regions: Self::Instruction,
35 ) -> Option<InstructionReply> {
36 let requests = open_regions
37 .into_iter()
38 .map(|open_region| {
39 let OpenRegion {
40 region_ident,
41 region_storage_path,
42 mut region_options,
43 region_wal_options,
44 skip_wal_replay,
45 reason,
46 requirements,
47 } = open_region;
48 let region_id = RegionId::new(region_ident.table_id, region_ident.region_number);
49 info!(
50 "Received open region instruction, region_id: {region_id}, reason: {reason:?}"
51 );
52 if let Err(err) =
53 serialize_wal_options(&mut region_options, region_id, ®ion_wal_options)
54 {
55 return Err(format!(
56 "Failed to serialize WAL options for region {region_id}: {err:?}"
57 ));
58 }
59 let request = RegionOpenRequest {
60 engine: region_ident.engine,
61 table_dir: table_dir(®ion_storage_path, region_id.table_id()),
62 path_type: PathType::Bare,
63 options: region_options,
64 skip_wal_replay,
65 checkpoint: None,
66 requirements,
67 };
68 Ok((region_id, request))
69 })
70 .collect::<Result<Vec<_>, _>>();
71 let requests = match requests {
72 Ok(requests) => requests,
73 Err(error) => {
74 return Some(InstructionReply::OpenRegions(SimpleReply {
75 result: false,
76 error: Some(InstructionError::legacy_internal_retryable(error)),
77 }));
78 }
79 };
80
81 let result = ctx
82 .region_server
83 .handle_batch_open_requests(self.open_region_parallelism, requests, false)
84 .await;
85 let success = result.is_ok();
86 let error = result.as_ref().map_err(InstructionError::from_error).err();
87
88 Some(InstructionReply::OpenRegions(SimpleReply {
89 result: success,
90 error,
91 }))
92 }
93}
94
95#[cfg(test)]
96mod tests {
97 use std::assert_matches;
98 use std::collections::HashMap;
99 use std::sync::Arc;
100
101 use common_error::ext::{BoxedError, RetryHint};
102 use common_error::status_code::StatusCode;
103 use common_meta::RegionIdent;
104 use common_meta::heartbeat::handler::{HandleControl, HeartbeatResponseHandler};
105 use common_meta::heartbeat::mailbox::MessageMeta;
106 use common_meta::instruction::{Instruction, OpenRegion};
107 use common_meta::kv_backend::memory::MemoryKvBackend;
108 use mito2::config::MitoConfig;
109 use mito2::engine::MITO_ENGINE_NAME;
110 use mito2::error::ManifestDeltaNotFoundSnafu;
111 use mito2::test_util::{CreateRequestBuilder, TestEnv};
112 use object_store::{Error as ObjectStoreError, ErrorKind};
113 use snafu::IntoError;
114 use store_api::path_utils::table_dir;
115 use store_api::region_request::{RegionCloseRequest, RegionRequest, RegionRequirements};
116 use store_api::storage::RegionId;
117
118 use super::OpenRegionsHandler;
119 use crate::error::{self, HandleRegionRequestSnafu};
120 use crate::heartbeat::handler::tests::HeartbeatResponseTestEnv;
121 use crate::heartbeat::handler::{
122 HandlerContext, InstructionHandler, RegionHeartbeatResponseHandler,
123 };
124 use crate::tests::{MockRegionEngine, mock_region_server};
125
126 fn open_regions_instruction(
127 region_ids: impl IntoIterator<Item = RegionId>,
128 storage_path: &str,
129 ) -> Instruction {
130 let region_idents = region_ids
131 .into_iter()
132 .map(|region_id| {
133 OpenRegion::new(
134 RegionIdent {
135 datanode_id: 0,
136 table_id: region_id.table_id(),
137 region_number: region_id.region_number(),
138 engine: MITO_ENGINE_NAME.to_string(),
139 },
140 storage_path,
141 HashMap::new(),
142 HashMap::new(),
143 false,
144 None,
145 RegionRequirements::empty(),
146 )
147 })
148 .collect();
149
150 Instruction::OpenRegions(region_idents)
151 }
152
153 #[tokio::test]
154 async fn test_open_regions() {
155 common_telemetry::init_default_ut_logging();
156
157 let mut region_server = mock_region_server();
158 let kv_backend = Arc::new(MemoryKvBackend::new());
159 let heartbeat_handler =
160 RegionHeartbeatResponseHandler::new(region_server.clone(), kv_backend);
161 let mut engine_env = TestEnv::with_prefix("open-regions").await;
162 let engine = engine_env.create_engine(MitoConfig::default()).await;
163 region_server.register_engine(Arc::new(engine.clone()));
164 let region_id = RegionId::new(1024, 1);
165 let region_id1 = RegionId::new(1024, 2);
166 let storage_path = "test";
167 let builder = CreateRequestBuilder::new();
168 let mut create_req = builder.build();
169 create_req.table_dir = table_dir(storage_path, region_id.table_id());
170 region_server
171 .handle_request(region_id, RegionRequest::Create(create_req))
172 .await
173 .unwrap();
174 let mut create_req1 = builder.build();
175 create_req1.table_dir = table_dir(storage_path, region_id1.table_id());
176 region_server
177 .handle_request(region_id1, RegionRequest::Create(create_req1))
178 .await
179 .unwrap();
180 region_server
181 .handle_request(
182 region_id,
183 RegionRequest::Close(RegionCloseRequest::default()),
184 )
185 .await
186 .unwrap();
187 region_server
188 .handle_request(
189 region_id,
190 RegionRequest::Close(RegionCloseRequest::default()),
191 )
192 .await
193 .unwrap();
194
195 let meta = MessageMeta::new_test(1, "test", "dn-1", "me-0");
196 let instruction = open_regions_instruction([region_id, region_id1], storage_path);
197 let mut heartbeat_env = HeartbeatResponseTestEnv::new();
198 let mut ctx = heartbeat_env.create_handler_ctx((meta, Default::default(), instruction));
199 let control = heartbeat_handler.handle(&mut ctx).await.unwrap();
200 assert_matches!(control, HandleControl::Continue);
201 let (_, reply) = heartbeat_env.receiver.recv().await.unwrap();
202
203 let reply = reply.expect_open_regions_reply();
204 assert!(reply.result);
205 assert!(reply.error.is_none());
206
207 assert!(engine.is_region_exists(region_id));
208 assert!(engine.is_region_exists(region_id1));
209 }
210
211 #[tokio::test]
212 async fn test_open_regions_preserves_manifest_delta_not_found_retry_hint() {
213 let retryable_region = RegionId::new(1024, 1);
214 let non_retryable_region = RegionId::new(1024, 2);
215 let (engine, _) = MockRegionEngine::with_mock_fn(
216 MITO_ENGINE_NAME,
217 Box::new(move |region_id, _request| {
218 if region_id == retryable_region {
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 return Err(HandleRegionRequestSnafu { region_id }
228 .into_error(BoxedError::new(manifest_error)));
229 }
230
231 error::RegionNotFoundSnafu { region_id }.fail()
232 }),
233 );
234 let mut region_server = mock_region_server();
235 region_server.register_engine(engine);
236 let ctx = HandlerContext::new_for_test(region_server, Arc::new(MemoryKvBackend::new()));
237 let Instruction::OpenRegions(open_regions) =
238 open_regions_instruction([retryable_region, non_retryable_region], "test")
239 else {
240 unreachable!()
241 };
242
243 let reply = OpenRegionsHandler {
245 open_region_parallelism: 1,
246 }
247 .handle(&ctx, open_regions)
248 .await
249 .unwrap()
250 .expect_open_regions_reply();
251
252 assert!(!reply.result);
253 let error = reply.error.unwrap();
254 assert_eq!(StatusCode::StorageUnavailable, error.code);
255 assert_eq!(RetryHint::Retryable, error.retry_hint);
256 assert!(error.message.contains("00000000000000000001.json"));
257 }
258}