1use std::any::Any;
18use std::collections::HashMap;
19use std::fmt::{Debug, Display};
20use std::sync::{Arc, Mutex};
21
22use api::greptime_proto::v1::meta::{GrantedRegion as PbGrantedRegion, RegionRole as PbRegionRole};
23use api::region::RegionResponse;
24use async_trait::async_trait;
25use common_error::ext::BoxedError;
26use common_recordbatch::adapter::RegionQueryStatCounters;
27use common_recordbatch::{EmptyRecordBatchStream, QueryMemoryTracker, SendableRecordBatchStream};
28use common_time::Timestamp;
29use datafusion_physical_plan::metrics::ExecutionPlanMetricsSet;
30use datafusion_physical_plan::{DisplayAs, DisplayFormatType, PhysicalExpr};
31use datatypes::schema::SchemaRef;
32use futures::future::join_all;
33use serde::{Deserialize, Serialize};
34use tokio::sync::Semaphore;
35
36use crate::logstore::entry;
37use crate::metadata::RegionMetadataRef;
38use crate::region_request::{
39 BatchRegionDdlRequest, RegionCatchupRequest, RegionOpenRequest, RegionRequest,
40};
41use crate::storage::{FileId, RegionId, ScanRequest, SequenceNumber};
42
43#[derive(Debug, PartialEq, Eq, Clone, Copy)]
45pub enum SettableRegionRoleState {
46 Follower,
47 DowngradingLeader,
48 Leader,
50 StagingLeader,
52}
53
54impl Display for SettableRegionRoleState {
55 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
56 match self {
57 SettableRegionRoleState::Follower => write!(f, "Follower"),
58 SettableRegionRoleState::DowngradingLeader => write!(f, "Leader(Downgrading)"),
59 SettableRegionRoleState::Leader => write!(f, "Leader"),
60 SettableRegionRoleState::StagingLeader => write!(f, "Leader(Staging)"),
61 }
62 }
63}
64
65impl From<SettableRegionRoleState> for RegionRole {
66 fn from(value: SettableRegionRoleState) -> Self {
67 match value {
68 SettableRegionRoleState::Follower => RegionRole::Follower,
69 SettableRegionRoleState::DowngradingLeader => RegionRole::DowngradingLeader,
70 SettableRegionRoleState::Leader => RegionRole::Leader,
71 SettableRegionRoleState::StagingLeader => RegionRole::StagingLeader,
72 }
73 }
74}
75
76#[derive(Debug, PartialEq, Eq)]
78pub struct SetRegionRoleStateRequest {
79 region_id: RegionId,
80 region_role_state: SettableRegionRoleState,
81}
82
83#[derive(Debug, PartialEq, Eq)]
85pub enum SetRegionRoleStateSuccess {
86 File,
87 Mito {
88 last_entry_id: entry::Id,
89 },
90 Metric {
91 last_entry_id: entry::Id,
92 metadata_last_entry_id: entry::Id,
93 },
94}
95
96impl SetRegionRoleStateSuccess {
97 pub fn file() -> Self {
99 Self::File
100 }
101
102 pub fn mito(last_entry_id: entry::Id) -> Self {
104 SetRegionRoleStateSuccess::Mito { last_entry_id }
105 }
106
107 pub fn metric(last_entry_id: entry::Id, metadata_last_entry_id: entry::Id) -> Self {
109 SetRegionRoleStateSuccess::Metric {
110 last_entry_id,
111 metadata_last_entry_id,
112 }
113 }
114}
115
116impl SetRegionRoleStateSuccess {
117 pub fn last_entry_id(&self) -> Option<entry::Id> {
119 match self {
120 SetRegionRoleStateSuccess::File => None,
121 SetRegionRoleStateSuccess::Mito { last_entry_id } => Some(*last_entry_id),
122 SetRegionRoleStateSuccess::Metric { last_entry_id, .. } => Some(*last_entry_id),
123 }
124 }
125
126 pub fn metadata_last_entry_id(&self) -> Option<entry::Id> {
128 match self {
129 SetRegionRoleStateSuccess::File => None,
130 SetRegionRoleStateSuccess::Mito { .. } => None,
131 SetRegionRoleStateSuccess::Metric {
132 metadata_last_entry_id,
133 ..
134 } => Some(*metadata_last_entry_id),
135 }
136 }
137}
138
139#[derive(Debug)]
141pub enum SetRegionRoleStateResponse {
142 Success(SetRegionRoleStateSuccess),
143 NotFound,
144 InvalidTransition(BoxedError),
145}
146
147impl SetRegionRoleStateResponse {
148 pub fn success(success: SetRegionRoleStateSuccess) -> Self {
150 Self::Success(success)
151 }
152
153 pub fn invalid_transition(error: BoxedError) -> Self {
155 Self::InvalidTransition(error)
156 }
157
158 pub fn is_not_found(&self) -> bool {
160 matches!(self, SetRegionRoleStateResponse::NotFound)
161 }
162
163 pub fn is_invalid_transition(&self) -> bool {
165 matches!(self, SetRegionRoleStateResponse::InvalidTransition(_))
166 }
167}
168
169#[derive(Debug, Clone, PartialEq, Eq)]
170pub struct GrantedRegion {
171 pub region_id: RegionId,
172 pub region_role: RegionRole,
173 pub extensions: HashMap<String, Vec<u8>>,
174}
175
176impl GrantedRegion {
177 pub fn new(region_id: RegionId, region_role: RegionRole) -> Self {
178 Self {
179 region_id,
180 region_role,
181 extensions: HashMap::new(),
182 }
183 }
184}
185
186impl From<GrantedRegion> for PbGrantedRegion {
187 fn from(value: GrantedRegion) -> Self {
188 PbGrantedRegion {
189 region_id: value.region_id.as_u64(),
190 role: PbRegionRole::from(value.region_role).into(),
191 extensions: value.extensions,
192 }
193 }
194}
195
196impl From<PbGrantedRegion> for GrantedRegion {
197 fn from(value: PbGrantedRegion) -> Self {
198 GrantedRegion {
199 region_id: RegionId::from_u64(value.region_id),
200 region_role: value.role().into(),
201 extensions: value.extensions,
202 }
203 }
204}
205
206#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
209pub enum RegionRole {
210 Follower,
212 Leader,
214 StagingLeader,
219 DowngradingLeader,
223}
224
225impl Display for RegionRole {
226 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
227 match self {
228 RegionRole::Follower => write!(f, "Follower"),
229 RegionRole::Leader => write!(f, "Leader"),
230 RegionRole::StagingLeader => write!(f, "Leader(Staging)"),
231 RegionRole::DowngradingLeader => write!(f, "Leader(Downgrading)"),
232 }
233 }
234}
235
236impl RegionRole {
237 pub fn writable(&self) -> bool {
238 matches!(self, RegionRole::Leader | RegionRole::StagingLeader)
239 }
240}
241
242impl From<RegionRole> for PbRegionRole {
243 fn from(value: RegionRole) -> Self {
244 match value {
245 RegionRole::Follower => PbRegionRole::Follower,
246 RegionRole::Leader => PbRegionRole::Leader,
247 RegionRole::StagingLeader => PbRegionRole::StagingLeader,
248 RegionRole::DowngradingLeader => PbRegionRole::DowngradingLeader,
249 }
250 }
251}
252
253impl From<PbRegionRole> for RegionRole {
254 fn from(value: PbRegionRole) -> Self {
255 match value {
256 PbRegionRole::Leader => RegionRole::Leader,
257 PbRegionRole::StagingLeader => RegionRole::StagingLeader,
258 PbRegionRole::Follower => RegionRole::Follower,
259 PbRegionRole::DowngradingLeader => RegionRole::DowngradingLeader,
260 }
261 }
262}
263
264#[derive(Debug)]
266pub enum ScannerPartitioning {
267 Unknown(usize),
269}
270
271impl ScannerPartitioning {
272 pub fn num_partitions(&self) -> usize {
274 match self {
275 ScannerPartitioning::Unknown(num_partitions) => *num_partitions,
276 }
277 }
278}
279
280#[derive(Debug, Clone, Copy, PartialEq, Eq)]
282pub struct PartitionRange {
283 pub start: Timestamp,
285 pub end: Timestamp,
287 pub num_rows: usize,
289 pub identifier: usize,
291}
292
293#[derive(Debug, Default)]
295pub struct ScannerProperties {
296 pub partitions: Vec<Vec<PartitionRange>>,
301
302 append_mode: bool,
304
305 total_rows: usize,
308
309 pub distinguish_partition_range: bool,
311
312 target_partitions: usize,
314
315 logical_region: bool,
317
318 query_load_region_id: Option<RegionId>,
320 query_stat_counters: Option<RegionQueryStatCounters>,
322}
323
324impl ScannerProperties {
325 pub fn with_append_mode(mut self, append_mode: bool) -> Self {
327 self.append_mode = append_mode;
328 self
329 }
330
331 pub fn with_total_rows(mut self, total_rows: usize) -> Self {
333 self.total_rows = total_rows;
334 self
335 }
336
337 pub fn new(partitions: Vec<Vec<PartitionRange>>, append_mode: bool, total_rows: usize) -> Self {
339 Self {
340 partitions,
341 append_mode,
342 total_rows,
343 distinguish_partition_range: false,
344 target_partitions: 0,
345 logical_region: false,
346 query_load_region_id: None,
347 query_stat_counters: None,
348 }
349 }
350
351 pub fn prepare(&mut self, request: PrepareRequest) {
353 if let Some(ranges) = request.ranges {
354 self.partitions = ranges;
355 }
356 if let Some(distinguish_partition_range) = request.distinguish_partition_range {
357 self.distinguish_partition_range = distinguish_partition_range;
358 }
359 if let Some(target_partitions) = request.target_partitions {
360 self.target_partitions = target_partitions;
361 }
362 }
363
364 pub fn num_partitions(&self) -> usize {
366 self.partitions.len()
367 }
368
369 pub fn append_mode(&self) -> bool {
370 self.append_mode
371 }
372
373 pub fn total_rows(&self) -> usize {
374 self.total_rows
375 }
376
377 pub fn is_logical_region(&self) -> bool {
379 self.logical_region
380 }
381
382 pub fn target_partitions(&self) -> usize {
384 if self.target_partitions == 0 {
385 self.num_partitions()
386 } else {
387 self.target_partitions
388 }
389 }
390
391 pub fn set_logical_region(&mut self, logical_region: bool) {
393 self.logical_region = logical_region;
394 }
395
396 pub fn query_load_region_id(&self) -> Option<RegionId> {
398 self.query_load_region_id
399 }
400
401 pub fn query_stat_counters(&self) -> Option<RegionQueryStatCounters> {
403 self.query_stat_counters.clone()
404 }
405
406 pub fn set_query_load_region_id(&mut self, region_id: RegionId) {
408 self.query_load_region_id = Some(region_id);
409 }
410
411 pub fn set_query_stat_counters(&mut self, counters: RegionQueryStatCounters) {
413 self.query_stat_counters = Some(counters);
414 }
415}
416
417#[derive(Default)]
419pub struct PrepareRequest {
420 pub ranges: Option<Vec<Vec<PartitionRange>>>,
422 pub distinguish_partition_range: Option<bool>,
424 pub target_partitions: Option<usize>,
426}
427
428impl PrepareRequest {
429 pub fn with_ranges(mut self, ranges: Vec<Vec<PartitionRange>>) -> Self {
431 self.ranges = Some(ranges);
432 self
433 }
434
435 pub fn with_distinguish_partition_range(mut self, distinguish_partition_range: bool) -> Self {
437 self.distinguish_partition_range = Some(distinguish_partition_range);
438 self
439 }
440
441 pub fn with_target_partitions(mut self, target_partitions: usize) -> Self {
443 self.target_partitions = Some(target_partitions);
444 self
445 }
446}
447
448#[derive(Clone, Default)]
450pub struct QueryScanContext {
451 pub explain_verbose: bool,
453}
454
455pub trait RegionScanner: Debug + DisplayAs + Send {
460 fn name(&self) -> &str;
461
462 fn properties(&self) -> &ScannerProperties;
464
465 fn schema(&self) -> SchemaRef;
467
468 fn metadata(&self) -> RegionMetadataRef;
470
471 fn prepare(&mut self, request: PrepareRequest) -> Result<(), BoxedError>;
475
476 fn scan_partition(
481 &self,
482 ctx: &QueryScanContext,
483 metrics_set: &ExecutionPlanMetricsSet,
484 partition: usize,
485 ) -> Result<SendableRecordBatchStream, BoxedError>;
486
487 fn has_predicate_without_region(&self) -> bool;
489
490 fn add_dyn_filter_to_predicate(
495 &mut self,
496 filter_exprs: Vec<Arc<dyn PhysicalExpr>>,
497 ) -> Vec<bool>;
498
499 fn set_logical_region(&mut self, logical_region: bool);
501
502 fn set_query_load_region_id(&mut self, region_id: RegionId);
504
505 fn snapshot_sequence(&self) -> Option<SequenceNumber> {
506 None
507 }
508}
509
510pub type RegionScannerRef = Box<dyn RegionScanner>;
511
512pub type BatchResponses = Vec<(RegionId, Result<RegionResponse, BoxedError>)>;
513
514#[derive(Debug, Deserialize, Serialize, Default)]
516pub struct RegionStatistic {
517 #[serde(default)]
522 pub num_rows: u64,
523 pub memtable_size: u64,
525 pub wal_size: u64,
527 pub manifest_size: u64,
529 pub sst_size: u64,
533 pub sst_num: u64,
537 #[serde(default)]
541 pub index_size: u64,
542 #[serde(default)]
544 pub manifest: RegionManifestInfo,
545 #[serde(default)]
546 pub written_bytes: u64,
548 #[serde(default)]
552 pub query_cpu_time: u64,
553 #[serde(default)]
555 pub query_scanned_bytes: u64,
556 #[serde(default)]
560 pub data_topic_latest_entry_id: u64,
561 #[serde(default)]
562 pub metadata_topic_latest_entry_id: u64,
563}
564
565#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
567pub enum RegionManifestInfo {
568 Mito {
569 manifest_version: u64,
570 flushed_entry_id: u64,
571 file_removed_cnt: u64,
573 },
574 Metric {
575 data_manifest_version: u64,
576 data_flushed_entry_id: u64,
577 metadata_manifest_version: u64,
578 metadata_flushed_entry_id: u64,
579 },
580}
581
582impl RegionManifestInfo {
583 pub fn mito(manifest_version: u64, flushed_entry_id: u64, file_removal_rate: u64) -> Self {
585 Self::Mito {
586 manifest_version,
587 flushed_entry_id,
588 file_removed_cnt: file_removal_rate,
589 }
590 }
591
592 pub fn metric(
594 data_manifest_version: u64,
595 data_flushed_entry_id: u64,
596 metadata_manifest_version: u64,
597 metadata_flushed_entry_id: u64,
598 ) -> Self {
599 Self::Metric {
600 data_manifest_version,
601 data_flushed_entry_id,
602 metadata_manifest_version,
603 metadata_flushed_entry_id,
604 }
605 }
606
607 pub fn is_mito(&self) -> bool {
609 matches!(self, RegionManifestInfo::Mito { .. })
610 }
611
612 pub fn is_metric(&self) -> bool {
614 matches!(self, RegionManifestInfo::Metric { .. })
615 }
616
617 pub fn data_flushed_entry_id(&self) -> u64 {
619 match self {
620 RegionManifestInfo::Mito {
621 flushed_entry_id, ..
622 } => *flushed_entry_id,
623 RegionManifestInfo::Metric {
624 data_flushed_entry_id,
625 ..
626 } => *data_flushed_entry_id,
627 }
628 }
629
630 pub fn data_manifest_version(&self) -> u64 {
632 match self {
633 RegionManifestInfo::Mito {
634 manifest_version, ..
635 } => *manifest_version,
636 RegionManifestInfo::Metric {
637 data_manifest_version,
638 ..
639 } => *data_manifest_version,
640 }
641 }
642
643 pub fn metadata_manifest_version(&self) -> Option<u64> {
645 match self {
646 RegionManifestInfo::Mito { .. } => None,
647 RegionManifestInfo::Metric {
648 metadata_manifest_version,
649 ..
650 } => Some(*metadata_manifest_version),
651 }
652 }
653
654 pub fn metadata_flushed_entry_id(&self) -> Option<u64> {
656 match self {
657 RegionManifestInfo::Mito { .. } => None,
658 RegionManifestInfo::Metric {
659 metadata_flushed_entry_id,
660 ..
661 } => Some(*metadata_flushed_entry_id),
662 }
663 }
664
665 pub fn encode_list(manifest_infos: &[(RegionId, Self)]) -> serde_json::Result<Vec<u8>> {
667 serde_json::to_vec(manifest_infos)
668 }
669
670 pub fn decode_list(value: &[u8]) -> serde_json::Result<Vec<(RegionId, Self)>> {
672 serde_json::from_slice(value)
673 }
674}
675
676impl Default for RegionManifestInfo {
677 fn default() -> Self {
678 Self::Mito {
679 manifest_version: 0,
680 flushed_entry_id: 0,
681 file_removed_cnt: 0,
682 }
683 }
684}
685
686impl RegionStatistic {
687 pub fn deserialize_from_slice(value: &[u8]) -> Option<RegionStatistic> {
691 serde_json::from_slice(value).ok()
692 }
693
694 pub fn serialize_to_vec(&self) -> Option<Vec<u8>> {
698 serde_json::to_vec(self).ok()
699 }
700}
701
702impl RegionStatistic {
703 pub fn estimated_disk_size(&self) -> u64 {
705 self.wal_size + self.sst_size + self.manifest_size + self.index_size
706 }
707}
708
709#[cfg(test)]
710mod tests {
711 use serde_json::json;
712
713 use super::*;
714
715 #[test]
716 fn region_statistic_deserializes_without_query_stats() {
717 let statistic: RegionStatistic = serde_json::from_value(json!({
718 "num_rows": 1,
719 "memtable_size": 2,
720 "wal_size": 3,
721 "manifest_size": 4,
722 "sst_size": 5,
723 "sst_num": 6,
724 "index_size": 7,
725 "manifest": {
726 "Mito": {
727 "manifest_version": 8,
728 "flushed_entry_id": 9,
729 "file_removed_cnt": 10
730 }
731 },
732 "written_bytes": 11
733 }))
734 .unwrap();
735
736 assert_eq!(statistic.query_cpu_time, 0);
737 assert_eq!(statistic.query_scanned_bytes, 0);
738 }
739}
740
741#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
743pub enum SyncRegionFromRequest {
744 FromManifest(RegionManifestInfo),
747 FromRegion {
753 source_region_id: RegionId,
755 parallelism: usize,
757 },
758}
759
760impl From<RegionManifestInfo> for SyncRegionFromRequest {
761 fn from(manifest_info: RegionManifestInfo) -> Self {
762 SyncRegionFromRequest::FromManifest(manifest_info)
763 }
764}
765
766impl SyncRegionFromRequest {
767 pub fn from_manifest(manifest_info: RegionManifestInfo) -> Self {
769 SyncRegionFromRequest::FromManifest(manifest_info)
770 }
771
772 pub fn from_region(source_region_id: RegionId, parallelism: usize) -> Self {
774 SyncRegionFromRequest::FromRegion {
775 source_region_id,
776 parallelism,
777 }
778 }
779
780 pub fn is_from_manifest(&self) -> bool {
782 matches!(self, SyncRegionFromRequest::FromManifest { .. })
783 }
784
785 pub fn into_region_manifest_info(self) -> Option<RegionManifestInfo> {
789 match self {
790 SyncRegionFromRequest::FromManifest(manifest_info) => Some(manifest_info),
791 SyncRegionFromRequest::FromRegion { .. } => None,
792 }
793 }
794}
795
796#[derive(Debug)]
798pub enum SyncRegionFromResponse {
799 NotSupported,
800 Mito {
801 synced: bool,
803 },
804 Metric {
805 metadata_synced: bool,
807 data_synced: bool,
809 new_opened_logical_region_ids: Vec<RegionId>,
812 },
813}
814
815impl SyncRegionFromResponse {
816 pub fn is_data_synced(&self) -> bool {
818 match self {
819 SyncRegionFromResponse::NotSupported => false,
820 SyncRegionFromResponse::Mito { synced } => *synced,
821 SyncRegionFromResponse::Metric { data_synced, .. } => *data_synced,
822 }
823 }
824
825 pub fn is_mito(&self) -> bool {
827 matches!(self, SyncRegionFromResponse::Mito { .. })
828 }
829
830 pub fn is_metric(&self) -> bool {
832 matches!(self, SyncRegionFromResponse::Metric { .. })
833 }
834
835 pub fn new_opened_logical_region_ids(self) -> Option<Vec<RegionId>> {
837 match self {
838 SyncRegionFromResponse::Metric {
839 new_opened_logical_region_ids,
840 ..
841 } => Some(new_opened_logical_region_ids),
842 _ => None,
843 }
844 }
845}
846
847#[derive(Debug, Clone)]
849pub struct RemapManifestsRequest {
850 pub region_id: RegionId,
852 pub input_regions: Vec<RegionId>,
854 pub region_mapping: HashMap<RegionId, Vec<RegionId>>,
856 pub new_partition_exprs: HashMap<RegionId, String>,
858}
859
860#[derive(Debug, Clone)]
862pub struct RemapManifestsResponse {
863 pub manifest_paths: HashMap<RegionId, String>,
868}
869
870#[derive(Debug, Clone)]
872pub struct MitoCopyRegionFromRequest {
873 pub source_region_id: RegionId,
875 pub parallelism: usize,
877}
878
879#[derive(Debug, Clone)]
880pub struct MitoCopyRegionFromResponse {
881 pub copied_file_ids: Vec<FileId>,
883}
884
885#[async_trait]
886pub trait RegionEngine: Send + Sync {
887 fn name(&self) -> &str;
889
890 async fn handle_batch_open_requests(
892 &self,
893 parallelism: usize,
894 requests: Vec<(RegionId, RegionOpenRequest)>,
895 ) -> Result<BatchResponses, BoxedError> {
896 let semaphore = Arc::new(Semaphore::new(parallelism));
897 let mut tasks = Vec::with_capacity(requests.len());
898
899 for (region_id, request) in requests {
900 let semaphore_moved = semaphore.clone();
901
902 tasks.push(async move {
903 let _permit = semaphore_moved.acquire().await.unwrap();
905 let result = self
906 .handle_request(region_id, RegionRequest::Open(request))
907 .await;
908 (region_id, result)
909 });
910 }
911
912 Ok(join_all(tasks).await)
913 }
914
915 async fn handle_batch_catchup_requests(
916 &self,
917 parallelism: usize,
918 requests: Vec<(RegionId, RegionCatchupRequest)>,
919 ) -> Result<BatchResponses, BoxedError> {
920 let semaphore = Arc::new(Semaphore::new(parallelism));
921 let mut tasks = Vec::with_capacity(requests.len());
922
923 for (region_id, request) in requests {
924 let semaphore_moved = semaphore.clone();
925
926 tasks.push(async move {
927 let _permit = semaphore_moved.acquire().await.unwrap();
929 let result = self
930 .handle_request(region_id, RegionRequest::Catchup(request))
931 .await;
932 (region_id, result)
933 });
934 }
935
936 Ok(join_all(tasks).await)
937 }
938
939 async fn handle_batch_ddl_requests(
940 &self,
941 request: BatchRegionDdlRequest,
942 ) -> Result<RegionResponse, BoxedError> {
943 let requests = request.into_region_requests();
944
945 let mut affected_rows = 0;
946 let mut extensions = HashMap::new();
947
948 for (region_id, request) in requests {
949 let result = self.handle_request(region_id, request).await?;
950 affected_rows += result.affected_rows;
951 extensions.extend(result.extensions);
952 }
953
954 Ok(RegionResponse {
955 affected_rows,
956 extensions,
957 metadata: Vec::new(),
958 })
959 }
960
961 async fn handle_request(
963 &self,
964 region_id: RegionId,
965 request: RegionRequest,
966 ) -> Result<RegionResponse, BoxedError>;
967
968 async fn get_committed_sequence(
970 &self,
971 region_id: RegionId,
972 ) -> Result<SequenceNumber, BoxedError>;
973
974 async fn handle_query(
976 &self,
977 region_id: RegionId,
978 request: ScanRequest,
979 ) -> Result<RegionScannerRef, BoxedError>;
980
981 fn query_memory_tracker(&self) -> Option<QueryMemoryTracker> {
983 None
984 }
985
986 async fn get_metadata(&self, region_id: RegionId) -> Result<RegionMetadataRef, BoxedError>;
988
989 fn region_statistic(&self, region_id: RegionId) -> Option<RegionStatistic>;
991
992 async fn stop(&self) -> Result<(), BoxedError>;
994
995 fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<(), BoxedError>;
1001
1002 async fn sync_region(
1004 &self,
1005 region_id: RegionId,
1006 request: SyncRegionFromRequest,
1007 ) -> Result<SyncRegionFromResponse, BoxedError>;
1008
1009 async fn remap_manifests(
1011 &self,
1012 request: RemapManifestsRequest,
1013 ) -> Result<RemapManifestsResponse, BoxedError>;
1014
1015 async fn set_region_role_state_gracefully(
1019 &self,
1020 region_id: RegionId,
1021 region_role_state: SettableRegionRoleState,
1022 ) -> Result<SetRegionRoleStateResponse, BoxedError>;
1023
1024 fn role(&self, region_id: RegionId) -> Option<RegionRole>;
1028
1029 fn as_any(&self) -> &dyn Any;
1030}
1031
1032pub type RegionEngineRef = Arc<dyn RegionEngine>;
1033
1034pub struct SinglePartitionScanner {
1036 stream: Mutex<Option<SendableRecordBatchStream>>,
1037 schema: SchemaRef,
1038 properties: ScannerProperties,
1039 metadata: RegionMetadataRef,
1040 snapshot_sequence: Option<SequenceNumber>,
1041}
1042
1043impl SinglePartitionScanner {
1044 pub fn new(
1046 stream: SendableRecordBatchStream,
1047 append_mode: bool,
1048 metadata: RegionMetadataRef,
1049 snapshot_sequence: Option<SequenceNumber>,
1050 ) -> Self {
1051 let schema = stream.schema();
1052 Self {
1053 stream: Mutex::new(Some(stream)),
1054 schema,
1055 properties: ScannerProperties::default().with_append_mode(append_mode),
1056 metadata,
1057 snapshot_sequence,
1058 }
1059 }
1060}
1061
1062impl Debug for SinglePartitionScanner {
1063 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1064 write!(f, "SinglePartitionScanner: <SendableRecordBatchStream>")
1065 }
1066}
1067
1068impl RegionScanner for SinglePartitionScanner {
1069 fn name(&self) -> &str {
1070 "SinglePartition"
1071 }
1072
1073 fn properties(&self) -> &ScannerProperties {
1074 &self.properties
1075 }
1076
1077 fn schema(&self) -> SchemaRef {
1078 self.schema.clone()
1079 }
1080
1081 fn prepare(&mut self, request: PrepareRequest) -> Result<(), BoxedError> {
1082 self.properties.prepare(request);
1083 Ok(())
1084 }
1085
1086 fn scan_partition(
1087 &self,
1088 _ctx: &QueryScanContext,
1089 _metrics_set: &ExecutionPlanMetricsSet,
1090 _partition: usize,
1091 ) -> Result<SendableRecordBatchStream, BoxedError> {
1092 let mut stream = self.stream.lock().unwrap();
1093 let result = stream
1094 .take()
1095 .or_else(|| Some(Box::pin(EmptyRecordBatchStream::new(self.schema.clone()))));
1096 Ok(result.unwrap())
1097 }
1098
1099 fn has_predicate_without_region(&self) -> bool {
1100 false
1101 }
1102
1103 fn add_dyn_filter_to_predicate(
1104 &mut self,
1105 filter_exprs: Vec<Arc<dyn datafusion_physical_plan::PhysicalExpr>>,
1106 ) -> Vec<bool> {
1107 vec![false; filter_exprs.len()]
1108 }
1109
1110 fn metadata(&self) -> RegionMetadataRef {
1111 self.metadata.clone()
1112 }
1113
1114 fn set_logical_region(&mut self, logical_region: bool) {
1115 self.properties.set_logical_region(logical_region);
1116 }
1117
1118 fn set_query_load_region_id(&mut self, region_id: RegionId) {
1119 self.properties.set_query_load_region_id(region_id);
1120 }
1121
1122 fn snapshot_sequence(&self) -> Option<SequenceNumber> {
1123 self.snapshot_sequence
1124 }
1125}
1126
1127impl DisplayAs for SinglePartitionScanner {
1128 fn fmt_as(&self, _t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
1129 write!(f, "{:?}", self)
1130 }
1131}