1use 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 ®ion_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(®ion_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 ®ion_name(region_id.table_id(), region_id.region_sequence()),
84 );
85 let manifest = FileRegionManifest::load(region_id, ®ion_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 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(®ion, &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}