Skip to main content

datanode/heartbeat/handler/
sync_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::{
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/// Handler for [SyncRegion] instruction.
24/// It syncs the region from a manifest or another region.
25#[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    /// Handles a batch of [SyncRegion] instructions.
33    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    /// Handles a single [SyncRegion] instruction.
53    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}