1use 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 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 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 regions: RwLock<HashMap<RegionId, FileRegionRef>>,
185
186 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 Ok(())
252 }
253
254 fn state(&self, region_id: RegionId) -> Option<RegionRole> {
255 if self.regions.read().unwrap().get(®ion_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 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 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(®ion_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(®ion, &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(®ion_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(®ion_id).cloned()
384 }
385
386 async fn exists(&self, region_id: RegionId) -> bool {
387 self.regions.read().unwrap().contains_key(®ion_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}