Skip to main content

datanode/heartbeat/handler/
open_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 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, &region_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(&region_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        // Serial execution makes the first error selection deterministic.
244        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}