Skip to main content

file_engine/
engine.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 std::any::Any;
16use std::collections::HashMap;
17use std::sync::{Arc, RwLock};
18
19use api::region::RegionResponse;
20use async_trait::async_trait;
21use common_catalog::consts::FILE_ENGINE;
22use common_datasource::object_store::LocalFileAccess;
23use common_error::ext::BoxedError;
24use common_recordbatch::SendableRecordBatchStream;
25use common_telemetry::{error, info};
26use object_store::ObjectStore;
27use snafu::{OptionExt, ensure};
28use store_api::metadata::RegionMetadataRef;
29use store_api::region_engine::{
30    RegionEngine, RegionRole, RegionScannerRef, RegionStatistic, RemapManifestsRequest,
31    RemapManifestsResponse, SetRegionRoleStateResponse, SetRegionRoleStateSuccess,
32    SettableRegionRoleState, SinglePartitionScanner, SyncRegionFromRequest, SyncRegionFromResponse,
33};
34use store_api::region_request::{
35    AffectedRows, RegionCloseRequest, RegionCreateRequest, RegionDropRequest, RegionOpenRequest,
36    RegionRequest, RegionRequirements,
37};
38use store_api::storage::{RegionId, ScanRequest, SequenceNumber};
39use tokio::sync::Mutex;
40
41use crate::config::EngineConfig;
42use crate::error::{
43    RegionNotFoundSnafu, Result as EngineResult, UnexpectedEngineSnafu, UnsupportedSnafu,
44};
45use crate::region::{FileRegion, FileRegionRef};
46
47pub struct FileRegionEngine {
48    inner: EngineInnerRef,
49}
50
51impl FileRegionEngine {
52    pub fn new(
53        _config: EngineConfig,
54        object_store: ObjectStore,
55        local_file_access: LocalFileAccess,
56    ) -> Self {
57        Self {
58            inner: Arc::new(EngineInner::new(object_store, local_file_access)),
59        }
60    }
61
62    async fn handle_query(
63        &self,
64        region_id: RegionId,
65        request: ScanRequest,
66    ) -> Result<SendableRecordBatchStream, BoxedError> {
67        self.inner
68            .get_region(region_id)
69            .await
70            .context(RegionNotFoundSnafu { region_id })
71            .map_err(BoxedError::new)?
72            .query(request, &self.inner.local_file_access)
73            .await
74            .map_err(BoxedError::new)
75    }
76}
77
78#[async_trait]
79impl RegionEngine for FileRegionEngine {
80    fn name(&self) -> &str {
81        FILE_ENGINE
82    }
83
84    async fn handle_request(
85        &self,
86        region_id: RegionId,
87        request: RegionRequest,
88    ) -> Result<RegionResponse, BoxedError> {
89        self.inner
90            .handle_request(region_id, request)
91            .await
92            .map_err(BoxedError::new)
93    }
94
95    async fn handle_query(
96        &self,
97        region_id: RegionId,
98        request: ScanRequest,
99    ) -> Result<RegionScannerRef, BoxedError> {
100        let stream = self.handle_query(region_id, request).await?;
101        let metadata = self.get_metadata(region_id).await?;
102        // We don't support enabling append mode for file engine.
103        let scanner = Box::new(SinglePartitionScanner::new(stream, false, metadata, None));
104        Ok(scanner)
105    }
106
107    async fn get_metadata(&self, region_id: RegionId) -> Result<RegionMetadataRef, BoxedError> {
108        self.inner
109            .get_region(region_id)
110            .await
111            .map(|r| r.metadata())
112            .context(RegionNotFoundSnafu { region_id })
113            .map_err(BoxedError::new)
114    }
115
116    async fn stop(&self) -> Result<(), BoxedError> {
117        self.inner.stop().await.map_err(BoxedError::new)
118    }
119
120    fn region_statistic(&self, _: RegionId) -> Option<RegionStatistic> {
121        None
122    }
123
124    async fn get_committed_sequence(&self, _: RegionId) -> Result<SequenceNumber, BoxedError> {
125        Ok(Default::default())
126    }
127
128    fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<(), BoxedError> {
129        self.inner
130            .set_region_role(region_id, role)
131            .map_err(BoxedError::new)
132    }
133
134    async fn set_region_role_state_gracefully(
135        &self,
136        region_id: RegionId,
137        _region_role_state: SettableRegionRoleState,
138    ) -> Result<SetRegionRoleStateResponse, BoxedError> {
139        let exists = self.inner.get_region(region_id).await.is_some();
140
141        if exists {
142            Ok(SetRegionRoleStateResponse::success(
143                SetRegionRoleStateSuccess::file(),
144            ))
145        } else {
146            Ok(SetRegionRoleStateResponse::NotFound)
147        }
148    }
149
150    async fn sync_region(
151        &self,
152        _region_id: RegionId,
153        _request: SyncRegionFromRequest,
154    ) -> Result<SyncRegionFromResponse, BoxedError> {
155        // File engine doesn't need to sync region manifest.
156        Ok(SyncRegionFromResponse::NotSupported)
157    }
158
159    async fn remap_manifests(
160        &self,
161        _request: RemapManifestsRequest,
162    ) -> Result<RemapManifestsResponse, BoxedError> {
163        Err(BoxedError::new(
164            UnsupportedSnafu {
165                operation: "remap_manifests",
166            }
167            .build(),
168        ))
169    }
170
171    fn role(&self, region_id: RegionId) -> Option<RegionRole> {
172        self.inner.state(region_id)
173    }
174
175    fn as_any(&self) -> &dyn Any {
176        self
177    }
178}
179
180struct EngineInner {
181    /// All regions opened by the engine.
182    ///
183    /// Writing to `regions` should also hold the `region_mutex`.
184    regions: RwLock<HashMap<RegionId, FileRegionRef>>,
185
186    /// Region mutex is used to protect the operations such as creating/opening/closing
187    /// a region, to avoid things like opening the same region simultaneously.
188    region_mutex: Mutex<()>,
189
190    object_store: ObjectStore,
191
192    local_file_access: LocalFileAccess,
193}
194
195type EngineInnerRef = Arc<EngineInner>;
196
197fn ensure_region_requirements(
198    requirements: RegionRequirements,
199    object_store: &ObjectStore,
200) -> EngineResult<()> {
201    if !requirements.object_storage {
202        return Ok(());
203    }
204
205    ensure!(
206        object_store::util::is_object_storage(object_store),
207        UnsupportedSnafu {
208            operation: "open region with object storage requirement on non-object storage"
209        }
210    );
211
212    Ok(())
213}
214
215impl EngineInner {
216    fn new(object_store: ObjectStore, local_file_access: LocalFileAccess) -> Self {
217        Self {
218            regions: RwLock::new(HashMap::new()),
219            region_mutex: Mutex::new(()),
220            object_store,
221            local_file_access,
222        }
223    }
224
225    async fn handle_request(
226        &self,
227        region_id: RegionId,
228        request: RegionRequest,
229    ) -> EngineResult<RegionResponse> {
230        let result = match request {
231            RegionRequest::Create(req) => self.handle_create(region_id, req).await,
232            RegionRequest::Drop(req) => self.handle_drop(region_id, req).await,
233            RegionRequest::Open(req) => self.handle_open(region_id, req).await,
234            RegionRequest::Close(req) => self.handle_close(region_id, req).await,
235            _ => UnsupportedSnafu {
236                operation: request.to_string(),
237            }
238            .fail(),
239        };
240        result.map(RegionResponse::new)
241    }
242
243    async fn stop(&self) -> EngineResult<()> {
244        let _lock = self.region_mutex.lock().await;
245        self.regions.write().unwrap().clear();
246        Ok(())
247    }
248
249    fn set_region_role(&self, _region_id: RegionId, _region_role: RegionRole) -> EngineResult<()> {
250        // TODO(zhongzc): Improve the semantics and implementation of this API.
251        Ok(())
252    }
253
254    fn state(&self, region_id: RegionId) -> Option<RegionRole> {
255        if self.regions.read().unwrap().get(&region_id).is_some() {
256            Some(RegionRole::Leader)
257        } else {
258            None
259        }
260    }
261}
262
263impl EngineInner {
264    async fn handle_create(
265        &self,
266        region_id: RegionId,
267        request: RegionCreateRequest,
268    ) -> EngineResult<AffectedRows> {
269        ensure!(
270            request.engine == FILE_ENGINE,
271            UnexpectedEngineSnafu {
272                engine: request.engine
273            }
274        );
275
276        if self.exists(region_id).await {
277            return Ok(0);
278        }
279
280        info!("Try to create region, region_id: {}", region_id);
281
282        let _lock = self.region_mutex.lock().await;
283        // Check again after acquiring the lock
284        if self.exists(region_id).await {
285            return Ok(0);
286        }
287
288        ensure_region_requirements(request.requirements, &self.object_store)?;
289
290        let res = FileRegion::create(region_id, request, &self.object_store).await;
291        let region = res.inspect_err(|err| {
292            error!(
293                err;
294                "Failed to create region, region_id: {}",
295                region_id
296            );
297        })?;
298        self.regions.write().unwrap().insert(region_id, region);
299
300        info!("A new region is created, region_id: {}", region_id);
301        Ok(0)
302    }
303
304    async fn handle_open(
305        &self,
306        region_id: RegionId,
307        request: RegionOpenRequest,
308    ) -> EngineResult<AffectedRows> {
309        if self.exists(region_id).await {
310            return Ok(0);
311        }
312
313        info!("Try to open region, region_id: {}", region_id);
314
315        let _lock = self.region_mutex.lock().await;
316        // Check again after acquiring the lock
317        if self.exists(region_id).await {
318            return Ok(0);
319        }
320
321        ensure_region_requirements(request.requirements, &self.object_store)?;
322
323        let res = FileRegion::open(region_id, request, &self.object_store).await;
324        let region = res.inspect_err(|err| {
325            error!(
326                err;
327                "Failed to open region, region_id: {}",
328                region_id
329            );
330        })?;
331        self.regions.write().unwrap().insert(region_id, region);
332
333        info!("Region opened, region_id: {}", region_id);
334        Ok(0)
335    }
336
337    async fn handle_close(
338        &self,
339        region_id: RegionId,
340        _request: RegionCloseRequest,
341    ) -> EngineResult<AffectedRows> {
342        let _lock = self.region_mutex.lock().await;
343
344        let mut regions = self.regions.write().unwrap();
345        if regions.remove(&region_id).is_some() {
346            info!("Region closed, region_id: {}", region_id);
347        }
348
349        Ok(0)
350    }
351
352    async fn handle_drop(
353        &self,
354        region_id: RegionId,
355        _request: RegionDropRequest,
356    ) -> EngineResult<AffectedRows> {
357        if !self.exists(region_id).await {
358            return RegionNotFoundSnafu { region_id }.fail();
359        }
360
361        info!("Try to drop region, region_id: {}", region_id);
362
363        let _lock = self.region_mutex.lock().await;
364
365        let region = self.get_region(region_id).await;
366        if let Some(region) = region {
367            let res = FileRegion::drop(&region, &self.object_store).await;
368            res.inspect_err(|err| {
369                error!(
370                    err;
371                    "Failed to drop region, region_id: {}",
372                    region_id
373                );
374            })?;
375        }
376        let _ = self.regions.write().unwrap().remove(&region_id);
377
378        info!("Region dropped, region_id: {}", region_id);
379        Ok(0)
380    }
381
382    async fn get_region(&self, region_id: RegionId) -> Option<FileRegionRef> {
383        self.regions.read().unwrap().get(&region_id).cloned()
384    }
385
386    async fn exists(&self, region_id: RegionId) -> bool {
387        self.regions.read().unwrap().contains_key(&region_id)
388    }
389}
390
391#[cfg(test)]
392mod tests {
393    use object_store::services::{Fs, S3};
394
395    use super::*;
396    use crate::error::Error;
397
398    fn build_fs_object_store() -> ObjectStore {
399        ObjectStore::new(Fs::default().root("/tmp"))
400            .unwrap()
401            .finish()
402    }
403
404    fn build_s3_object_store() -> ObjectStore {
405        ObjectStore::new(
406            S3::default()
407                .bucket("test-bucket")
408                .region("us-east-1")
409                .disable_ec2_metadata(),
410        )
411        .unwrap()
412        .finish()
413    }
414
415    #[test]
416    fn test_empty_region_requirements_are_supported() {
417        ensure_region_requirements(RegionRequirements::empty(), &build_fs_object_store()).unwrap();
418    }
419
420    #[test]
421    fn test_object_storage_region_requirement_rejects_fs_object_store() {
422        let err = ensure_region_requirements(
423            RegionRequirements::object_storage(),
424            &build_fs_object_store(),
425        )
426        .unwrap_err();
427
428        assert!(matches!(err, Error::Unsupported { .. }));
429    }
430
431    #[test]
432    fn test_object_storage_region_requirement_accepts_s3_object_store() {
433        ensure_region_requirements(
434            RegionRequirements::object_storage(),
435            &build_s3_object_store(),
436        )
437        .unwrap();
438    }
439}