Skip to main content

file_engine/
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 std::collections::HashMap;
16use std::sync::Arc;
17
18use common_datasource::file_format::Format;
19use object_store::ObjectStore;
20use store_api::metadata::RegionMetadataRef;
21use store_api::path_utils::region_name;
22use store_api::region_request::{RegionCreateRequest, RegionOpenRequest};
23use store_api::storage::RegionId;
24
25use crate::FileOptions;
26use crate::error::Result;
27use crate::manifest::FileRegionManifest;
28
29#[derive(Debug)]
30pub struct FileRegion {
31    pub(crate) region_dir: String,
32    pub(crate) file_options: FileOptions,
33    pub(crate) url: String,
34    pub(crate) format: Format,
35    pub(crate) options: HashMap<String, String>,
36    pub(crate) metadata: RegionMetadataRef,
37}
38
39pub type FileRegionRef = Arc<FileRegion>;
40
41impl FileRegion {
42    pub async fn create(
43        region_id: RegionId,
44        request: RegionCreateRequest,
45        object_store: &ObjectStore,
46    ) -> Result<FileRegionRef> {
47        let manifest = FileRegionManifest {
48            region_id,
49            column_metadatas: request.column_metadatas.clone(),
50            primary_key: request.primary_key.clone(),
51            options: request.options,
52        };
53
54        let region_dir = object_store::util::join_dir(
55            &request.table_dir,
56            &region_name(region_id.table_id(), region_id.region_sequence()),
57        );
58        let url = manifest.url()?;
59        let file_options = manifest.file_options()?;
60        let format = manifest.format()?;
61        let options = manifest.options.clone();
62        let metadata = manifest.metadata()?;
63
64        manifest.store(&region_dir, object_store).await?;
65
66        Ok(Arc::new(Self {
67            region_dir,
68            url,
69            file_options,
70            format,
71            options,
72            metadata,
73        }))
74    }
75
76    pub async fn open(
77        region_id: RegionId,
78        request: RegionOpenRequest,
79        object_store: &ObjectStore,
80    ) -> Result<FileRegionRef> {
81        let region_dir = object_store::util::join_dir(
82            &request.table_dir,
83            &region_name(region_id.table_id(), region_id.region_sequence()),
84        );
85        let manifest = FileRegionManifest::load(region_id, &region_dir, object_store).await?;
86
87        Ok(Arc::new(Self {
88            region_dir,
89            url: manifest.url()?,
90            file_options: manifest.file_options()?,
91            format: manifest.format()?,
92            metadata: manifest.metadata()?,
93            options: manifest.options,
94        }))
95    }
96
97    pub async fn drop(&self, object_store: &ObjectStore) -> Result<()> {
98        FileRegionManifest::delete(self.metadata.region_id, &self.region_dir, object_store).await
99    }
100
101    pub fn metadata(&self) -> RegionMetadataRef {
102        self.metadata.clone()
103    }
104}
105
106#[cfg(test)]
107mod tests {
108    use std::assert_matches;
109
110    use common_datasource::object_store::LocalFileAccess;
111    use common_error::ext::{ErrorExt, RetryHint};
112    use common_error::status_code::StatusCode;
113    use store_api::region_request::PathType;
114    use store_api::storage::ScanRequest;
115
116    use super::*;
117    use crate::error::Error;
118    use crate::test_util::{new_test_column_metadata, new_test_object_store, new_test_options};
119
120    #[tokio::test]
121    async fn test_create_region() {
122        let (_dir, object_store) = new_test_object_store("test_create_region");
123
124        let request = RegionCreateRequest {
125            engine: "file".to_string(),
126            column_metadatas: new_test_column_metadata(),
127            primary_key: vec![1],
128            options: new_test_options(),
129            table_dir: "create_region_dir/".to_string(),
130            path_type: PathType::Bare,
131            partition_expr_json: Some("".to_string()),
132            requirements: Default::default(),
133        };
134        let region_id = RegionId::new(1, 0);
135
136        let region = FileRegion::create(region_id, request.clone(), &object_store)
137            .await
138            .unwrap();
139
140        assert_eq!(region.region_dir, "create_region_dir/1_0000000000/");
141        assert_eq!(region.url, "test");
142        assert_eq!(region.file_options.files, vec!["1.csv"]);
143        assert_matches!(region.format, Format::Csv { .. });
144        assert_eq!(region.options, new_test_options());
145        assert_eq!(region.metadata.region_id, region_id);
146        assert_eq!(region.metadata.primary_key, vec![1]);
147
148        assert!(
149            object_store
150                .exists("create_region_dir/1_0000000000/manifest/_file_manifest")
151                .await
152                .unwrap()
153        );
154
155        // Object exists, should fail
156        let err = FileRegion::create(region_id, request, &object_store)
157            .await
158            .unwrap_err();
159        assert_matches!(err, Error::ManifestExists { .. });
160    }
161
162    #[tokio::test]
163    async fn test_persisted_local_region_rejected_when_disabled() {
164        let (_dir, object_store) = new_test_object_store("test_disabled_local_region");
165        let request = RegionCreateRequest {
166            engine: "file".to_string(),
167            column_metadatas: new_test_column_metadata(),
168            primary_key: vec![1],
169            options: new_test_options(),
170            table_dir: "disabled_local_region/".to_string(),
171            path_type: PathType::Bare,
172            partition_expr_json: Some("".to_string()),
173            requirements: Default::default(),
174        };
175        let region = FileRegion::create(RegionId::new(1, 0), request, &object_store)
176            .await
177            .unwrap();
178
179        let error = match region
180            .query(ScanRequest::default(), &LocalFileAccess::Disabled)
181            .await
182        {
183            Ok(_) => panic!("local file query must be rejected"),
184            Err(error) => error,
185        };
186        assert_matches!(
187            &error,
188            Error::BuildBackend {
189                source: common_datasource::error::Error::LocalFileAccessDisabled { .. },
190                ..
191            }
192        );
193        assert_eq!(error.status_code(), StatusCode::InvalidArguments);
194        assert_eq!(error.retry_hint(), RetryHint::NonRetryable);
195    }
196
197    #[tokio::test]
198    async fn test_open_region() {
199        let (_dir, object_store) = new_test_object_store("test_open_region");
200
201        let region_dir = "open_region_dir/".to_string();
202        let request = RegionCreateRequest {
203            engine: "file".to_string(),
204            column_metadatas: new_test_column_metadata(),
205            primary_key: vec![1],
206            options: new_test_options(),
207            table_dir: region_dir.clone(),
208            path_type: PathType::Bare,
209            partition_expr_json: Some("".to_string()),
210            requirements: Default::default(),
211        };
212        let region_id = RegionId::new(1, 0);
213
214        let _ = FileRegion::create(region_id, request.clone(), &object_store)
215            .await
216            .unwrap();
217
218        let request = RegionOpenRequest {
219            engine: "file".to_string(),
220            table_dir: region_dir,
221            path_type: PathType::Bare,
222            options: HashMap::default(),
223            skip_wal_replay: false,
224            checkpoint: None,
225            requirements: Default::default(),
226        };
227
228        let region = FileRegion::open(region_id, request, &object_store)
229            .await
230            .unwrap();
231
232        assert_eq!(region.region_dir, "open_region_dir/1_0000000000/");
233        assert_eq!(region.url, "test");
234        assert_eq!(region.file_options.files, vec!["1.csv"]);
235        assert_matches!(region.format, Format::Csv { .. });
236        assert_eq!(region.options, new_test_options());
237        assert_eq!(region.metadata.region_id, region_id);
238        assert_eq!(region.metadata.primary_key, vec![1]);
239    }
240
241    #[tokio::test]
242    async fn test_drop_region() {
243        let (_dir, object_store) = new_test_object_store("test_drop_region");
244
245        let region_dir = "drop_region_dir/".to_string();
246        let request = RegionCreateRequest {
247            engine: "file".to_string(),
248            column_metadatas: new_test_column_metadata(),
249            primary_key: vec![1],
250            options: new_test_options(),
251            table_dir: region_dir.clone(),
252            path_type: PathType::Bare,
253            partition_expr_json: Some("".to_string()),
254            requirements: Default::default(),
255        };
256        let region_id = RegionId::new(1, 0);
257
258        let region = FileRegion::create(region_id, request.clone(), &object_store)
259            .await
260            .unwrap();
261
262        assert!(
263            object_store
264                .exists("drop_region_dir/1_0000000000/manifest/_file_manifest")
265                .await
266                .unwrap()
267        );
268
269        FileRegion::drop(&region, &object_store).await.unwrap();
270        assert!(
271            !object_store
272                .exists("drop_region_dir/1_0000000000/manifest/_file_manifest")
273                .await
274                .unwrap()
275        );
276
277        let request = RegionOpenRequest {
278            engine: "file".to_string(),
279            table_dir: region_dir,
280            path_type: PathType::Bare,
281            options: HashMap::default(),
282            skip_wal_replay: false,
283            checkpoint: None,
284            requirements: Default::default(),
285        };
286        let err = FileRegion::open(region_id, request, &object_store)
287            .await
288            .unwrap_err();
289        assert_matches!(err, Error::LoadRegionManifest { .. });
290    }
291}