1mod catalog;
16mod registrations;
17mod remote_dyn_filter;
18
19use std::collections::HashMap;
20use std::fmt::Debug;
21use std::ops::Deref;
22use std::pin::Pin;
23use std::sync::atomic::{AtomicBool, Ordering};
24use std::sync::{Arc, RwLock};
25use std::task::{Context, Poll};
26use std::time::Duration;
27
28use api::region::RegionResponse;
29use api::v1::meta::TopicStat;
30use api::v1::region::sync_request::ManifestInfo;
31use api::v1::region::{
32 ListMetadataRequest, RegionResponse as RegionResponseV1, SyncRequest, region_request,
33};
34use api::v1::{ResponseHeader, Status};
35use arrow_flight::{FlightData, Ticket};
36use async_trait::async_trait;
37use bytes::Bytes;
38use common_error::ext::{BoxedError, ErrorExt};
39use common_error::status_code::StatusCode;
40use common_meta::datanode::TopicStatsReporter;
41use common_query::OutputData;
42use common_query::request::QueryRequest;
43use common_recordbatch::adapter::RecordBatchMetrics;
44use common_recordbatch::{OrderOption, RecordBatch, RecordBatchStream, SendableRecordBatchStream};
45use common_runtime::Runtime;
46use common_telemetry::tracing::{self, info_span};
47use common_telemetry::tracing_context::{FutureExt, TracingContext};
48use common_telemetry::{debug, error, info, warn};
49use dashmap::DashMap;
50use datafusion::datasource::TableProvider;
51use datafusion_common::tree_node::TreeNode;
52use datatypes::schema::SchemaRef;
53use either::Either;
54use futures_util::future::try_join_all;
55use futures_util::{Stream, StreamExt};
56use metric_engine::engine::MetricEngine;
57use mito2::engine::{MITO_ENGINE_NAME, MitoEngine};
58use prost::Message;
59use query::QueryEngineRef;
60pub use query::dummy_catalog::{
61 DummyCatalogList, DummyTableProviderFactory, TableProviderFactoryRef,
62};
63use query::options::should_collect_region_watermark_from_extensions;
64use serde_json;
65use servers::error::{
66 self as servers_error, ExecuteGrpcRequestSnafu, Result as ServerResult, SuspendedSnafu,
67};
68use servers::grpc::FlightCompression;
69use servers::grpc::flight::{
70 FlightCraft, FlightRecordBatchSource, FlightRecordBatchStream, FlightRecordBatchStreamInput,
71 TonicStream,
72};
73use servers::grpc::region_server::RegionServerHandler;
74use session::context::{
75 FLIGHT_METRICS_HEARTBEAT_INTERVAL, QueryContext, QueryContextBuilder, QueryContextRef,
76};
77use snafu::{OptionExt, ResultExt, ensure};
78use store_api::metric_engine_consts::{
79 FILE_ENGINE_NAME, LOGICAL_TABLE_METADATA_KEY, METRIC_ENGINE_NAME,
80};
81use store_api::region_engine::{
82 RegionEngineRef, RegionManifestInfo, RegionRole, RegionStatistic, RemapManifestsRequest,
83 RemapManifestsResponse, SetRegionRoleStateResponse, SettableRegionRoleState,
84 SyncRegionFromRequest,
85};
86use store_api::region_request::{
87 AffectedRows, BatchRegionDdlRequest, RegionCatchupRequest, RegionCloseRequest,
88 RegionOpenRequest, RegionRequest,
89};
90use store_api::storage::RegionId;
91use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc};
92use tokio::time::{self, timeout};
93use tonic::{Request, Response, Result as TonicResult};
94
95use crate::error::{
96 self, BuildRegionRequestsSnafu, ConcurrentQueryLimiterClosedSnafu,
97 ConcurrentQueryLimiterTimeoutSnafu, DataFusionSnafu, DecodeLogicalPlanSnafu,
98 ExecuteLogicalPlanSnafu, FindLogicalRegionsSnafu, GetRegionMetadataSnafu,
99 HandleBatchDdlRequestSnafu, HandleBatchOpenRequestSnafu, HandleRegionRequestSnafu,
100 NewPlanDecoderSnafu, RegionEngineNotFoundSnafu, RegionNotFoundSnafu, RegionNotReadySnafu,
101 Result, RuntimeJoinSnafu, SerializeJsonSnafu, StopRegionEngineSnafu, UnexpectedSnafu,
102 UnsupportedOutputSnafu,
103};
104use crate::event_listener::RegionServerEventListenerRef;
105use crate::query_stream::QueryRuntimeStream;
106use crate::region_server::catalog::{NameAwareCatalogList, NameAwareDataSourceInjectorBuilder};
107use crate::region_server::registrations::RemoteDynFilterRegistry;
108use crate::region_server::remote_dyn_filter::wrap_remote_dyn_filter_guarded_stream;
109
110const QUERY_RUNTIME_STREAM_BUFFER_SIZE: usize = 8;
111
112#[derive(Clone)]
113pub struct RegionServer {
114 inner: Arc<RegionServerInner>,
115 flight_compression: FlightCompression,
116 suspend: Arc<AtomicBool>,
117}
118
119pub struct RegionStat {
120 pub region_id: RegionId,
121 pub engine: String,
122 pub role: RegionRole,
123}
124
125impl RegionServer {
126 pub fn new(
127 query_engine: QueryEngineRef,
128 runtime: Runtime,
129 event_listener: RegionServerEventListenerRef,
130 flight_compression: FlightCompression,
131 ) -> Self {
132 Self::with_table_provider(
133 query_engine,
134 runtime,
135 event_listener,
136 Arc::new(DummyTableProviderFactory),
137 0,
138 Duration::from_millis(0),
139 flight_compression,
140 )
141 }
142
143 pub fn with_table_provider(
144 query_engine: QueryEngineRef,
145 runtime: Runtime,
146 event_listener: RegionServerEventListenerRef,
147 table_provider_factory: TableProviderFactoryRef,
148 max_concurrent_queries: usize,
149 concurrent_query_limiter_timeout: Duration,
150 flight_compression: FlightCompression,
151 ) -> Self {
152 Self {
153 inner: Arc::new(RegionServerInner::new(
154 query_engine,
155 runtime,
156 event_listener,
157 table_provider_factory,
158 RegionServerParallelism::from_opts(
159 max_concurrent_queries,
160 concurrent_query_limiter_timeout,
161 ),
162 )),
163 flight_compression,
164 suspend: Arc::new(AtomicBool::new(false)),
165 }
166 }
167
168 pub fn register_engine(&mut self, engine: RegionEngineRef) {
170 self.inner.register_engine(engine);
171 }
172
173 pub fn set_topic_stats_reporter(&mut self, topic_stats_reporter: Box<dyn TopicStatsReporter>) {
175 self.inner.set_topic_stats_reporter(topic_stats_reporter);
176 }
177
178 pub fn find_engine(&self, region_id: RegionId) -> Result<Option<RegionEngineRef>> {
180 match self.inner.get_engine(region_id, &RegionChange::None) {
181 Ok(CurrentEngine::Engine(engine)) => Ok(Some(engine)),
182 Ok(CurrentEngine::EarlyReturn(_)) => Ok(None),
183 Err(error::Error::RegionNotFound { .. }) => Ok(None),
184 Err(err) => Err(err),
185 }
186 }
187
188 pub fn mito_engine(&self) -> Option<MitoEngine> {
190 if let Some(mito) = self.inner.mito_engine.read().unwrap().clone() {
191 Some(mito)
192 } else {
193 self.inner
194 .engines
195 .read()
196 .unwrap()
197 .get(MITO_ENGINE_NAME)
198 .cloned()
199 .and_then(|e| {
200 let mito = e.as_any().downcast_ref::<MitoEngine>().cloned();
201 if mito.is_none() {
202 warn!("Mito engine not found in region server engines");
203 }
204 mito
205 })
206 }
207 }
208
209 #[tracing::instrument(skip_all)]
210 pub async fn handle_batch_open_requests(
211 &self,
212 parallelism: usize,
213 requests: Vec<(RegionId, RegionOpenRequest)>,
214 ignore_nonexistent_region: bool,
215 ) -> Result<Vec<RegionId>> {
216 self.inner
217 .handle_batch_open_requests(parallelism, requests, ignore_nonexistent_region)
218 .await
219 }
220
221 #[tracing::instrument(skip_all)]
222 pub async fn handle_batch_catchup_requests(
223 &self,
224 parallelism: usize,
225 requests: Vec<(RegionId, RegionCatchupRequest)>,
226 ) -> Result<Vec<(RegionId, std::result::Result<(), BoxedError>)>> {
227 self.inner
228 .handle_batch_catchup_requests(parallelism, requests)
229 .await
230 }
231
232 #[tracing::instrument(skip_all, fields(request_type = request.request_type()))]
233 pub async fn handle_request(
234 &self,
235 region_id: RegionId,
236 request: RegionRequest,
237 ) -> Result<RegionResponse> {
238 if RegionServerInner::is_ingest_request(&request) {
239 let inner = self.inner.clone();
240 let request_type = request.request_type();
241 return common_runtime::spawn_ingest(async move {
242 inner.handle_request(region_id, request).await
243 })
244 .await
245 .context(RuntimeJoinSnafu { request_type })?;
246 }
247
248 self.inner.handle_request(region_id, request).await
249 }
250
251 async fn table_provider(
253 &self,
254 region_id: RegionId,
255 ctx: Option<QueryContextRef>,
256 ) -> Result<Arc<dyn TableProvider>> {
257 let status = self
258 .inner
259 .region_map
260 .get(®ion_id)
261 .context(RegionNotFoundSnafu { region_id })?
262 .clone();
263 ensure!(
264 matches!(status, RegionEngineWithStatus::Ready(_)),
265 RegionNotReadySnafu { region_id }
266 );
267
268 let provider = self
269 .inner
270 .table_provider_factory
271 .create(region_id, status.into_engine(), ctx)
272 .await
273 .context(ExecuteLogicalPlanSnafu)?;
274
275 Ok(provider)
276 }
277
278 pub async fn handle_remote_read(
280 &self,
281 request: api::v1::region::QueryRequest,
282 query_ctx: QueryContextRef,
283 ) -> Result<SendableRecordBatchStream> {
284 self.handle_remote_read_inner(request, query_ctx).await
285 }
286
287 async fn handle_remote_read_inner(
288 &self,
289 request: api::v1::region::QueryRequest,
290 query_ctx: QueryContextRef,
291 ) -> Result<SendableRecordBatchStream> {
292 let permit = if let Some(p) = &self.inner.parallelism {
293 Some(p.acquire().await?)
294 } else {
295 None
296 };
297
298 let region_id = RegionId::from_u64(request.region_id);
299 let catalog_list = Arc::new(NameAwareCatalogList::new(
300 self.clone(),
301 region_id,
302 query_ctx.clone(),
303 ));
304
305 if query_ctx.explain_verbose() {
306 common_telemetry::info!("Handle remote read for region: {}", region_id);
307 }
308
309 let decoder = self
310 .inner
311 .query_engine
312 .engine_context(query_ctx.clone())
313 .new_plan_decoder()
314 .context(NewPlanDecoderSnafu)?;
315
316 let plan = decoder
317 .decode(Bytes::from(request.plan), catalog_list, false)
318 .await
319 .context(DecodeLogicalPlanSnafu)?;
320
321 let cleanup = self.register_initial_remote_dyn_filter_cleanup(&query_ctx, region_id);
322
323 let stream = self
324 .inner
325 .handle_read(
326 QueryRequest {
327 header: request.header,
328 region_id,
329 plan,
330 },
331 query_ctx.clone(),
332 )
333 .await?;
334
335 let stream = wrap_flow_region_watermark_stream(stream, region_id, &query_ctx);
336 let stream = if let Some(cleanup) = cleanup {
337 wrap_remote_dyn_filter_guarded_stream(stream, cleanup)
338 } else {
339 stream
340 };
341 Ok(maybe_guard_stream(stream, permit))
342 }
343
344 #[tracing::instrument(skip_all)]
345 pub async fn handle_read(&self, request: QueryRequest) -> Result<SendableRecordBatchStream> {
346 self.handle_read_inner(request).await
347 }
348
349 async fn handle_read_inner(&self, request: QueryRequest) -> Result<SendableRecordBatchStream> {
350 let permit = if let Some(p) = &self.inner.parallelism {
351 Some(p.acquire().await?)
352 } else {
353 None
354 };
355
356 let ctx = request.header.as_ref().map(|h| h.into());
357 let query_ctx = Arc::new(ctx.unwrap_or_else(|| QueryContextBuilder::default().build()));
358
359 let region_id = request.region_id;
360 let injector_builder = NameAwareDataSourceInjectorBuilder::from_plan(&request.plan)
361 .context(DataFusionSnafu)?;
362 let mut injector = injector_builder
363 .build(self, request.region_id, query_ctx.clone())
364 .await?;
365
366 let plan = request
367 .plan
368 .rewrite(&mut injector)
369 .context(DataFusionSnafu)?
370 .data;
371
372 let cleanup = self.register_initial_remote_dyn_filter_cleanup(&query_ctx, region_id);
373
374 let stream = self
375 .inner
376 .handle_read(QueryRequest { plan, ..request }, query_ctx.clone())
377 .await?;
378
379 let stream = wrap_flow_region_watermark_stream(stream, region_id, &query_ctx);
380 let stream = if let Some(cleanup) = cleanup {
381 wrap_remote_dyn_filter_guarded_stream(stream, cleanup)
382 } else {
383 stream
384 };
385 Ok(maybe_guard_stream(stream, permit))
386 }
387
388 pub fn reportable_regions(&self) -> Vec<RegionStat> {
392 self.inner
393 .region_map
394 .iter()
395 .filter_map(|e| {
396 let region_id = *e.key();
397 e.role(region_id).map(|role| RegionStat {
399 region_id,
400 engine: e.value().name().to_string(),
401 role,
402 })
403 })
404 .collect()
405 }
406
407 pub fn topic_stats(&self) -> Vec<TopicStat> {
409 let mut reporter = self.inner.topic_stats_reporter.write().unwrap();
410 let Some(reporter) = reporter.as_mut() else {
411 return vec![];
412 };
413 reporter
414 .reportable_topics()
415 .into_iter()
416 .map(|stat| TopicStat {
417 topic_name: stat.topic,
418 record_size: stat.record_size,
419 record_num: stat.record_num,
420 latest_entry_id: stat.latest_entry_id,
421 })
422 .collect()
423 }
424
425 pub fn is_region_leader(&self, region_id: RegionId) -> Option<bool> {
426 self.inner.region_map.get(®ion_id).and_then(|engine| {
427 engine.role(region_id).map(|role| match role {
428 RegionRole::Follower => false,
429 RegionRole::Leader => true,
430 RegionRole::StagingLeader => true,
431 RegionRole::DowngradingLeader => true,
432 })
433 })
434 }
435
436 pub fn set_region_role(&self, region_id: RegionId, role: RegionRole) -> Result<()> {
437 let engine = self
438 .inner
439 .region_map
440 .get(®ion_id)
441 .with_context(|| RegionNotFoundSnafu { region_id })?;
442 engine
443 .set_region_role(region_id, role)
444 .with_context(|_| HandleRegionRequestSnafu { region_id })
445 }
446
447 pub async fn set_region_role_state_gracefully(
457 &self,
458 region_id: RegionId,
459 state: SettableRegionRoleState,
460 ) -> Result<SetRegionRoleStateResponse> {
461 match self.inner.region_map.get(®ion_id) {
462 Some(engine) => Ok(engine
463 .set_region_role_state_gracefully(region_id, state)
464 .await
465 .with_context(|_| HandleRegionRequestSnafu { region_id })?),
466 None => Ok(SetRegionRoleStateResponse::NotFound),
467 }
468 }
469
470 pub fn runtime(&self) -> Runtime {
471 self.inner.runtime.clone()
472 }
473
474 pub fn region_statistic(&self, region_id: RegionId) -> Option<RegionStatistic> {
475 match self.inner.region_map.get(®ion_id) {
476 Some(e) => e.region_statistic(region_id),
477 None => None,
478 }
479 }
480
481 pub async fn stop(&self) -> Result<()> {
483 self.inner.stop().await
484 }
485
486 #[cfg(test)]
487 pub(crate) fn register_test_region(&self, region_id: RegionId, engine: RegionEngineRef) {
489 {
490 let mut engines = self.inner.engines.write().unwrap();
491 if !engines.contains_key(engine.name()) {
492 debug!("Registering test engine: {}", engine.name());
493 engines.insert(engine.name().to_string(), engine.clone());
494 }
495 }
496
497 self.inner
498 .region_map
499 .insert(region_id, RegionEngineWithStatus::Ready(engine));
500 }
501
502 async fn handle_batch_ddl_requests(
503 &self,
504 request: region_request::Body,
505 ) -> Result<RegionResponse> {
506 let batch_request = BatchRegionDdlRequest::try_from_request_body(request)
508 .context(BuildRegionRequestsSnafu)?
509 .unwrap();
510 let tracing_context = TracingContext::from_current_span();
511
512 let span = tracing_context.attach(info_span!("RegionServer::handle_batch_ddl_requests"));
513 self.inner
514 .handle_batch_request(batch_request)
515 .trace(span)
516 .await
517 }
518
519 async fn handle_requests_in_parallel(
520 &self,
521 request: region_request::Body,
522 ) -> Result<RegionResponse> {
523 let requests =
524 RegionRequest::try_from_request_body(request).context(BuildRegionRequestsSnafu)?;
525
526 match self.try_handle_metric_batch_puts(requests).await? {
529 Either::Left(response) => Ok(response),
530 Either::Right(requests) => {
531 let tracing_context = TracingContext::from_current_span();
533 let join_tasks =
534 requests
535 .into_iter()
536 .map(|(region_id, req): (RegionId, RegionRequest)| {
537 let self_to_move = self;
538 let span = tracing_context.attach(info_span!(
539 "RegionServer::handle_region_request",
540 region_id = region_id.to_string()
541 ));
542 async move {
543 self_to_move
544 .handle_request(region_id, req)
545 .trace(span)
546 .await
547 }
548 });
549
550 let results = try_join_all(join_tasks).await?;
551 let mut affected_rows = 0;
552 let mut extensions = HashMap::new();
553 for result in results {
554 affected_rows += result.affected_rows;
555 extensions.extend(result.extensions);
556 }
557
558 Ok(RegionResponse {
559 affected_rows,
560 extensions,
561 metadata: Vec::new(),
562 })
563 }
564 }
565 }
566
567 async fn handle_requests_in_serial(
568 &self,
569 request: region_request::Body,
570 ) -> Result<RegionResponse> {
571 let requests =
572 RegionRequest::try_from_request_body(request).context(BuildRegionRequestsSnafu)?;
573 let tracing_context = TracingContext::from_current_span();
574
575 let mut affected_rows = 0;
576 let mut extensions = HashMap::new();
577 for (region_id, req) in requests {
578 let span = tracing_context.attach(info_span!(
579 "RegionServer::handle_region_request",
580 region_id = region_id.to_string()
581 ));
582 let result = self.handle_request(region_id, req).trace(span).await?;
583
584 affected_rows += result.affected_rows;
585 extensions.extend(result.extensions);
586 }
587
588 Ok(RegionResponse {
589 affected_rows,
590 extensions,
591 metadata: Vec::new(),
592 })
593 }
594
595 async fn try_handle_metric_batch_puts(
602 &self,
603 requests: Vec<(RegionId, RegionRequest)>,
604 ) -> Result<Either<RegionResponse, Vec<(RegionId, RegionRequest)>>> {
605 if requests.is_empty() {
606 return Ok(Either::Right(requests));
607 }
608
609 if !matches!(requests[0].1, RegionRequest::Put(_)) {
611 return Ok(Either::Right(requests));
612 }
613 let first_region_id = requests[0].0;
614 let request_type = requests[0].1.request_type();
615
616 let engine = match self
620 .inner
621 .get_engine(first_region_id, &RegionChange::None)?
622 {
623 CurrentEngine::Engine(e) => e,
624 _ => return Ok(Either::Right(requests)),
625 };
626
627 if engine.name() != METRIC_ENGINE_NAME {
628 return Ok(Either::Right(requests));
629 }
630
631 let mut all_puts = true;
633 for (_, req) in &requests {
634 if !matches!(req, RegionRequest::Put(_)) {
635 all_puts = false;
636 break;
637 }
638 }
639
640 if !all_puts {
641 return Ok(Either::Right(requests));
642 }
643
644 let put_requests = requests.into_iter().map(|(region_id, req)| {
646 if let RegionRequest::Put(put) = req {
647 (region_id, put)
648 } else {
649 unreachable!("Already checked all are Put")
650 }
651 });
652
653 let metric_engine = engine
655 .as_any()
656 .downcast_ref::<MetricEngine>()
657 .context(UnexpectedSnafu {
658 violated: "Failed to downcast to MetricEngine",
659 })?
660 .clone();
661
662 let tracing_context = TracingContext::from_current_span();
663 let batch_size = put_requests.len();
664 let span = tracing_context.attach(info_span!(
665 "RegionServer::handle_metric_batch_puts",
666 batch_size = batch_size,
667 ));
668 let result = common_runtime::spawn_ingest(async move {
669 metric_engine
670 .put_regions_batch(put_requests)
671 .trace(span)
672 .await
673 })
674 .await
675 .context(RuntimeJoinSnafu { request_type })?
676 .map_err(BoxedError::new)
677 .context(HandleRegionRequestSnafu {
678 region_id: first_region_id,
679 });
680
681 match result {
682 Ok(total_affected) => {
683 crate::metrics::REGION_CHANGED_ROW_COUNT
684 .with_label_values(&[request_type])
685 .inc_by(total_affected as u64);
686 Ok(Either::Left(RegionResponse::new(total_affected)))
687 }
688 Err(err) => {
689 crate::metrics::REGION_SERVER_INSERT_FAIL_COUNT
690 .with_label_values(&[request_type])
691 .inc_by(batch_size as u64);
692 Err(err)
693 }
694 }
695 }
696
697 async fn handle_sync_region_request(&self, request: &SyncRequest) -> Result<RegionResponse> {
698 let region_id = RegionId::from_u64(request.region_id);
699 let manifest_info = request
700 .manifest_info
701 .context(error::MissingRequiredFieldSnafu {
702 name: "manifest_info",
703 })?;
704
705 let manifest_info = match manifest_info {
706 ManifestInfo::MitoManifestInfo(info) => {
707 RegionManifestInfo::mito(info.data_manifest_version, 0, 0)
708 }
709 ManifestInfo::MetricManifestInfo(info) => RegionManifestInfo::metric(
710 info.data_manifest_version,
711 0,
712 info.metadata_manifest_version,
713 0,
714 ),
715 };
716
717 let tracing_context = TracingContext::from_current_span();
718 let span = tracing_context.attach(info_span!("RegionServer::handle_sync_region_request"));
719
720 self.sync_region(
721 region_id,
722 SyncRegionFromRequest::from_manifest(manifest_info),
723 )
724 .trace(span)
725 .await
726 .map(|_| RegionResponse::new(AffectedRows::default()))
727 }
728
729 #[tracing::instrument(skip_all)]
734 async fn handle_list_metadata_request(
735 &self,
736 request: &ListMetadataRequest,
737 ) -> Result<RegionResponse> {
738 let mut region_metadatas = Vec::new();
739 for region_id in &request.region_ids {
741 let region_id = RegionId::from_u64(*region_id);
742 let Some(engine) = self.find_engine(region_id)? else {
744 region_metadatas.push(None);
745 continue;
746 };
747
748 match engine.get_metadata(region_id).await {
749 Ok(metadata) => region_metadatas.push(Some(metadata)),
750 Err(err) => {
751 if err.status_code() == StatusCode::RegionNotFound {
752 region_metadatas.push(None);
753 } else {
754 Err(err).with_context(|_| GetRegionMetadataSnafu {
755 engine: engine.name(),
756 region_id,
757 })?;
758 }
759 }
760 }
761 }
762
763 let json_result = serde_json::to_vec(®ion_metadatas).context(SerializeJsonSnafu)?;
765
766 let response = RegionResponse::from_metadata(json_result);
767
768 Ok(response)
769 }
770
771 pub async fn sync_region(
773 &self,
774 region_id: RegionId,
775 request: SyncRegionFromRequest,
776 ) -> Result<()> {
777 let engine = match self.inner.get_engine(region_id, &RegionChange::None)? {
778 CurrentEngine::Engine(engine) => engine,
779 _ => {
780 return UnexpectedSnafu {
781 violated: "unexpected EarlyReturn engine status for a ready region",
782 }
783 .fail();
784 }
785 };
786
787 self.inner
788 .handle_sync_region(&engine, region_id, request)
789 .await
790 }
791
792 pub async fn remap_manifests(
794 &self,
795 request: RemapManifestsRequest,
796 ) -> Result<RemapManifestsResponse> {
797 let region_id = request.region_id;
798 let engine = match self.inner.get_engine(region_id, &RegionChange::None)? {
799 CurrentEngine::Engine(engine) => engine,
800 _ => {
801 return UnexpectedSnafu {
802 violated: "unexpected EarlyReturn engine status for a ready region",
803 }
804 .fail();
805 }
806 };
807
808 engine
809 .remap_manifests(request)
810 .await
811 .with_context(|_| HandleRegionRequestSnafu { region_id })
812 }
813
814 fn is_suspended(&self) -> bool {
815 self.suspend.load(Ordering::Relaxed)
816 }
817
818 pub(crate) fn suspend_state(&self) -> Arc<AtomicBool> {
819 self.suspend.clone()
820 }
821}
822
823fn wrap_flow_region_watermark_stream(
824 stream: SendableRecordBatchStream,
825 region_id: RegionId,
826 query_ctx: &QueryContextRef,
827) -> SendableRecordBatchStream {
828 if should_collect_region_watermark_from_extensions(&query_ctx.extensions())
829 && let Some(seq) = query_ctx.get_snapshot(region_id.as_u64())
830 {
831 Box::pin(RegionWatermarkStream::new(stream, region_id, seq)) as SendableRecordBatchStream
832 } else {
833 stream
834 }
835}
836
837struct RegionWatermarkStream {
839 stream: SendableRecordBatchStream,
840 region_id: u64,
841 snapshot_seq: u64,
842 finished: bool,
843}
844
845impl RegionWatermarkStream {
846 fn new(stream: SendableRecordBatchStream, region_id: RegionId, snapshot_seq: u64) -> Self {
847 Self {
848 stream,
849 region_id: region_id.as_u64(),
850 snapshot_seq,
851 finished: false,
852 }
853 }
854
855 fn merged_metrics(&self, mut metrics: RecordBatchMetrics) -> RecordBatchMetrics {
856 if metrics
857 .region_watermarks
858 .iter()
859 .any(|entry| entry.region_id == self.region_id)
860 {
861 return metrics;
862 }
863
864 metrics
865 .region_watermarks
866 .push(common_recordbatch::adapter::RegionWatermarkEntry {
867 region_id: self.region_id,
868 watermark: Some(self.snapshot_seq),
869 });
870 metrics
871 }
872}
873
874impl RecordBatchStream for RegionWatermarkStream {
875 fn name(&self) -> &str {
876 self.stream.name()
877 }
878
879 fn schema(&self) -> datatypes::schema::SchemaRef {
880 self.stream.schema()
881 }
882
883 fn output_ordering(&self) -> Option<&[OrderOption]> {
884 self.stream.output_ordering()
885 }
886
887 fn metrics(&self) -> Option<RecordBatchMetrics> {
888 let base = self.stream.metrics();
889 if !self.finished {
890 return base;
891 }
892
893 Some(self.merged_metrics(base.unwrap_or_default()))
894 }
895}
896
897impl Stream for RegionWatermarkStream {
898 type Item = common_recordbatch::error::Result<RecordBatch>;
899
900 fn size_hint(&self) -> (usize, Option<usize>) {
901 self.stream.size_hint()
902 }
903
904 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
905 match Pin::new(&mut self.stream).poll_next(cx) {
906 Poll::Ready(None) => {
907 self.finished = true;
908 Poll::Ready(None)
909 }
910 other => other,
911 }
912 }
913}
914
915#[async_trait]
916impl RegionServerHandler for RegionServer {
917 async fn handle(&self, request: region_request::Body) -> ServerResult<RegionResponseV1> {
918 let failed_requests_cnt = crate::metrics::REGION_SERVER_REQUEST_FAILURE_COUNT
919 .with_label_values(&[request.as_ref()]);
920 let response = match &request {
921 region_request::Body::Creates(_)
922 | region_request::Body::Drops(_)
923 | region_request::Body::Alters(_) => self.handle_batch_ddl_requests(request).await,
924 region_request::Body::Inserts(_) | region_request::Body::Deletes(_) => {
925 self.handle_requests_in_parallel(request).await
926 }
927 region_request::Body::Sync(sync_request) => {
928 self.handle_sync_region_request(sync_request).await
929 }
930 region_request::Body::ListMetadata(list_metadata_request) => {
931 self.handle_list_metadata_request(list_metadata_request)
932 .await
933 }
934 region_request::Body::RemoteDynFilter(remote_dyn_filter_request) => {
935 self.handle_remote_dyn_filter_request(remote_dyn_filter_request)
936 .await
937 }
938 _ => self.handle_requests_in_serial(request).await,
939 }
940 .map_err(BoxedError::new)
941 .inspect_err(|_| {
942 failed_requests_cnt.inc();
943 })
944 .context(ExecuteGrpcRequestSnafu)?;
945
946 Ok(RegionResponseV1 {
947 header: Some(ResponseHeader {
948 status: Some(Status {
949 status_code: StatusCode::Success as _,
950 ..Default::default()
951 }),
952 }),
953 affected_rows: response.affected_rows as _,
954 extensions: response.extensions,
955 metadata: response.metadata,
956 })
957 }
958}
959
960#[async_trait]
961impl FlightCraft for RegionServer {
962 async fn do_get(
963 &self,
964 request: Request<Ticket>,
965 ) -> TonicResult<Response<TonicStream<FlightData>>> {
966 ensure!(!self.is_suspended(), SuspendedSnafu);
967
968 let ticket = request.into_inner().ticket;
969 let request = api::v1::region::QueryRequest::decode(ticket.as_ref())
970 .context(servers_error::InvalidFlightTicketSnafu)?;
971 let tracing_context = request
972 .header
973 .as_ref()
974 .map(|h| TracingContext::from_w3c(&h.tracing_context))
975 .unwrap_or_default();
976 let query_ctx = request
977 .header
978 .as_ref()
979 .map(|h| Arc::new(QueryContext::from(h)))
980 .unwrap_or(QueryContext::arc());
981
982 let region_server = self.clone();
983 let initializer_query_ctx = query_ctx.clone();
984 let initializer_tracing_context = tracing_context.clone();
985 let initializer = async move {
986 region_server
987 .handle_remote_read(request, initializer_query_ctx)
988 .trace(initializer_tracing_context.attach(info_span!("RegionServer::handle_read")))
989 .await
990 .map_err(Into::into)
991 };
992
993 let stream = Box::pin(FlightRecordBatchStream::new(
994 FlightRecordBatchStreamInput::initializer(async move {
995 initializer
996 .await
997 .map(FlightRecordBatchSource::RecordBatches)
998 }),
999 tracing_context,
1000 self.flight_compression,
1001 query_ctx,
1002 ));
1003 Ok(Response::new(stream))
1004 }
1005}
1006
1007#[derive(Clone)]
1008enum RegionEngineWithStatus {
1009 Registering(RegionEngineRef),
1011 Deregistering(RegionEngineRef),
1013 Ready(RegionEngineRef),
1015}
1016
1017impl RegionEngineWithStatus {
1018 pub fn into_engine(self) -> RegionEngineRef {
1020 match self {
1021 RegionEngineWithStatus::Registering(engine) => engine,
1022 RegionEngineWithStatus::Deregistering(engine) => engine,
1023 RegionEngineWithStatus::Ready(engine) => engine,
1024 }
1025 }
1026}
1027
1028impl Deref for RegionEngineWithStatus {
1029 type Target = RegionEngineRef;
1030
1031 fn deref(&self) -> &Self::Target {
1032 match self {
1033 RegionEngineWithStatus::Registering(engine) => engine,
1034 RegionEngineWithStatus::Deregistering(engine) => engine,
1035 RegionEngineWithStatus::Ready(engine) => engine,
1036 }
1037 }
1038}
1039
1040struct RegionServerInner {
1041 engines: RwLock<HashMap<String, RegionEngineRef>>,
1042 region_map: DashMap<RegionId, RegionEngineWithStatus>,
1043 query_engine: QueryEngineRef,
1044 runtime: Runtime,
1045 event_listener: RegionServerEventListenerRef,
1046 table_provider_factory: TableProviderFactoryRef,
1047 parallelism: Option<RegionServerParallelism>,
1050 topic_stats_reporter: RwLock<Option<Box<dyn TopicStatsReporter>>>,
1052 mito_engine: RwLock<Option<MitoEngine>>,
1056 initial_remote_dyn_filter_registrations: RemoteDynFilterRegistry,
1059}
1060
1061struct RegionServerParallelism {
1062 semaphore: Arc<Semaphore>,
1063 timeout: Duration,
1064}
1065
1066impl RegionServerParallelism {
1067 pub fn from_opts(
1068 max_concurrent_queries: usize,
1069 concurrent_query_limiter_timeout: Duration,
1070 ) -> Option<Self> {
1071 if max_concurrent_queries == 0 {
1072 return None;
1073 }
1074 Some(RegionServerParallelism {
1075 semaphore: Arc::new(Semaphore::new(max_concurrent_queries)),
1076 timeout: concurrent_query_limiter_timeout,
1077 })
1078 }
1079
1080 pub async fn acquire(&self) -> Result<OwnedSemaphorePermit> {
1081 timeout(self.timeout, self.semaphore.clone().acquire_owned())
1082 .await
1083 .context(ConcurrentQueryLimiterTimeoutSnafu)?
1084 .context(ConcurrentQueryLimiterClosedSnafu)
1085 }
1086}
1087
1088struct PermitGuardedStream {
1092 inner: SendableRecordBatchStream,
1093 _permit: OwnedSemaphorePermit,
1094}
1095
1096impl RecordBatchStream for PermitGuardedStream {
1097 fn name(&self) -> &str {
1098 self.inner.name()
1099 }
1100
1101 fn schema(&self) -> SchemaRef {
1102 self.inner.schema()
1103 }
1104
1105 fn output_ordering(&self) -> Option<&[OrderOption]> {
1106 self.inner.output_ordering()
1107 }
1108
1109 fn metrics(&self) -> Option<RecordBatchMetrics> {
1110 self.inner.metrics()
1111 }
1112}
1113
1114impl Stream for PermitGuardedStream {
1115 type Item = common_recordbatch::error::Result<RecordBatch>;
1116
1117 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
1118 self.inner.as_mut().poll_next(cx)
1119 }
1120}
1121
1122fn maybe_guard_stream(
1125 stream: SendableRecordBatchStream,
1126 permit: Option<OwnedSemaphorePermit>,
1127) -> SendableRecordBatchStream {
1128 match permit {
1129 Some(permit) => Box::pin(PermitGuardedStream {
1130 inner: stream,
1131 _permit: permit,
1132 }),
1133 None => stream,
1134 }
1135}
1136
1137enum CurrentEngine {
1138 Engine(RegionEngineRef),
1139 EarlyReturn(AffectedRows),
1140}
1141
1142impl Debug for CurrentEngine {
1143 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1144 match self {
1145 CurrentEngine::Engine(engine) => f
1146 .debug_struct("CurrentEngine")
1147 .field("engine", &engine.name())
1148 .finish(),
1149 CurrentEngine::EarlyReturn(rows) => f
1150 .debug_struct("CurrentEngine")
1151 .field("return", rows)
1152 .finish(),
1153 }
1154 }
1155}
1156
1157impl RegionServerInner {
1158 pub fn new(
1159 query_engine: QueryEngineRef,
1160 runtime: Runtime,
1161 event_listener: RegionServerEventListenerRef,
1162 table_provider_factory: TableProviderFactoryRef,
1163 parallelism: Option<RegionServerParallelism>,
1164 ) -> Self {
1165 Self {
1166 engines: RwLock::new(HashMap::new()),
1167 region_map: DashMap::new(),
1168 query_engine,
1169 runtime,
1170 event_listener,
1171 table_provider_factory,
1172 parallelism,
1173 topic_stats_reporter: RwLock::new(None),
1174 mito_engine: RwLock::new(None),
1175 initial_remote_dyn_filter_registrations: RemoteDynFilterRegistry::new(),
1176 }
1177 }
1178
1179 pub fn register_engine(&self, engine: RegionEngineRef) {
1180 let engine_name = engine.name();
1181 if engine_name == MITO_ENGINE_NAME
1182 && let Some(mito_engine) = engine.as_any().downcast_ref::<MitoEngine>()
1183 {
1184 *self.mito_engine.write().unwrap() = Some(mito_engine.clone());
1185 }
1186
1187 info!("Region Engine {engine_name} is registered");
1188 self.engines
1189 .write()
1190 .unwrap()
1191 .insert(engine_name.to_string(), engine);
1192 }
1193
1194 pub fn set_topic_stats_reporter(&self, topic_stats_reporter: Box<dyn TopicStatsReporter>) {
1195 info!("Set topic stats reporter");
1196 *self.topic_stats_reporter.write().unwrap() = Some(topic_stats_reporter);
1197 }
1198
1199 fn is_ingest_request(request: &RegionRequest) -> bool {
1200 matches!(
1201 request,
1202 RegionRequest::Put(_) | RegionRequest::Delete(_) | RegionRequest::BulkInserts(_)
1203 )
1204 }
1205
1206 fn get_engine(
1207 &self,
1208 region_id: RegionId,
1209 region_change: &RegionChange,
1210 ) -> Result<CurrentEngine> {
1211 let current_region_status = self.region_map.get(®ion_id);
1212
1213 let engine = match region_change {
1214 RegionChange::Register(attribute) => match current_region_status {
1215 Some(status) => match status.clone() {
1216 RegionEngineWithStatus::Registering(engine) => engine,
1217 RegionEngineWithStatus::Deregistering(_) => {
1218 return error::RegionBusySnafu { region_id }.fail();
1219 }
1220 RegionEngineWithStatus::Ready(_) => status.clone().into_engine(),
1221 },
1222 _ => self
1223 .engines
1224 .read()
1225 .unwrap()
1226 .get(attribute.engine())
1227 .with_context(|| RegionEngineNotFoundSnafu {
1228 name: attribute.engine(),
1229 })?
1230 .clone(),
1231 },
1232 RegionChange::OfflineCleanup(attribute) => match current_region_status {
1233 Some(status) => match status.clone() {
1234 RegionEngineWithStatus::Registering(_)
1235 | RegionEngineWithStatus::Deregistering(_)
1236 | RegionEngineWithStatus::Ready(_) => {
1237 return error::RegionBusySnafu { region_id }.fail();
1238 }
1239 },
1240 None => self
1241 .engines
1242 .read()
1243 .unwrap()
1244 .get(attribute.engine())
1245 .with_context(|| RegionEngineNotFoundSnafu {
1246 name: attribute.engine(),
1247 })?
1248 .clone(),
1249 },
1250 RegionChange::Deregisters => match current_region_status {
1251 Some(status) => match status.clone() {
1252 RegionEngineWithStatus::Registering(_) => {
1253 return error::RegionBusySnafu { region_id }.fail();
1254 }
1255 RegionEngineWithStatus::Deregistering(_) => {
1256 return Ok(CurrentEngine::EarlyReturn(0));
1257 }
1258 RegionEngineWithStatus::Ready(_) => status.clone().into_engine(),
1259 },
1260 None => return Ok(CurrentEngine::EarlyReturn(0)),
1261 },
1262 RegionChange::None | RegionChange::Catchup | RegionChange::Ingest => {
1263 match current_region_status {
1264 Some(status) => match status.clone() {
1265 RegionEngineWithStatus::Registering(_) => {
1266 return error::RegionNotReadySnafu { region_id }.fail();
1267 }
1268 RegionEngineWithStatus::Deregistering(_) => {
1269 return error::RegionNotFoundSnafu { region_id }.fail();
1270 }
1271 RegionEngineWithStatus::Ready(engine) => engine,
1272 },
1273 None => return error::RegionNotFoundSnafu { region_id }.fail(),
1274 }
1275 }
1276 };
1277
1278 Ok(CurrentEngine::Engine(engine))
1279 }
1280
1281 async fn handle_batch_open_requests_inner(
1282 &self,
1283 engine: RegionEngineRef,
1284 parallelism: usize,
1285 requests: Vec<(RegionId, RegionOpenRequest)>,
1286 ignore_nonexistent_region: bool,
1287 ) -> Result<Vec<RegionId>> {
1288 let region_changes = requests
1289 .iter()
1290 .map(|(region_id, open)| {
1291 let attribute = parse_region_attribute(&open.engine, &open.options)?;
1292 Ok((*region_id, RegionChange::Register(attribute)))
1293 })
1294 .collect::<Result<HashMap<_, _>>>()?;
1295
1296 for (®ion_id, region_change) in ®ion_changes {
1297 self.set_region_status_not_ready(region_id, &engine, region_change)
1298 }
1299
1300 let mut open_regions = Vec::with_capacity(requests.len());
1301 let mut errors = vec![];
1302 match engine
1303 .handle_batch_open_requests(parallelism, requests)
1304 .await
1305 .with_context(|_| HandleBatchOpenRequestSnafu)
1306 {
1307 Ok(results) => {
1308 for (region_id, result) in results {
1309 let region_change = ®ion_changes[®ion_id];
1310 match result {
1311 Ok(_) => {
1312 if let Err(e) = self
1313 .set_region_status_ready(region_id, engine.clone(), *region_change)
1314 .await
1315 {
1316 error!(e; "Failed to set region to ready: {}", region_id);
1317 errors.push(BoxedError::new(e));
1318 } else {
1319 open_regions.push(region_id)
1320 }
1321 }
1322 Err(e) => {
1323 self.unset_region_status(region_id, &engine, *region_change);
1324 if e.status_code() == StatusCode::RegionNotFound
1325 && ignore_nonexistent_region
1326 {
1327 warn!("Region {} not found, ignore it, source: {:?}", region_id, e);
1328 } else {
1329 error!(e; "Failed to open region: {}", region_id);
1330 errors.push(e);
1331 }
1332 }
1333 }
1334 }
1335 }
1336 Err(e) => {
1337 for (®ion_id, region_change) in ®ion_changes {
1338 self.unset_region_status(region_id, &engine, *region_change);
1339 }
1340 error!(e; "Failed to open batch regions");
1341 errors.push(BoxedError::new(e));
1342 }
1343 }
1344
1345 if !errors.is_empty() {
1346 return Err(errors.swap_remove(0)).context(HandleBatchOpenRequestSnafu);
1348 }
1349
1350 Ok(open_regions)
1351 }
1352
1353 pub async fn handle_batch_open_requests(
1354 &self,
1355 parallelism: usize,
1356 requests: Vec<(RegionId, RegionOpenRequest)>,
1357 ignore_nonexistent_region: bool,
1358 ) -> Result<Vec<RegionId>> {
1359 let mut engine_grouped_requests: HashMap<String, Vec<_>> =
1360 HashMap::with_capacity(requests.len());
1361 for (region_id, request) in requests {
1362 if let Some(requests) = engine_grouped_requests.get_mut(&request.engine) {
1363 requests.push((region_id, request));
1364 } else {
1365 engine_grouped_requests.insert(request.engine.clone(), vec![(region_id, request)]);
1366 }
1367 }
1368
1369 let mut results = Vec::with_capacity(engine_grouped_requests.keys().len());
1370 for (engine, requests) in engine_grouped_requests {
1371 let engine = self
1372 .engines
1373 .read()
1374 .unwrap()
1375 .get(&engine)
1376 .with_context(|| RegionEngineNotFoundSnafu { name: &engine })?
1377 .clone();
1378 results.push(
1379 self.handle_batch_open_requests_inner(
1380 engine,
1381 parallelism,
1382 requests,
1383 ignore_nonexistent_region,
1384 )
1385 .await,
1386 )
1387 }
1388
1389 Ok(results
1390 .into_iter()
1391 .collect::<Result<Vec<_>>>()?
1392 .into_iter()
1393 .flatten()
1394 .collect::<Vec<_>>())
1395 }
1396
1397 pub async fn handle_batch_catchup_requests_inner(
1398 &self,
1399 engine: RegionEngineRef,
1400 parallelism: usize,
1401 requests: Vec<(RegionId, RegionCatchupRequest)>,
1402 ) -> Result<Vec<(RegionId, std::result::Result<(), BoxedError>)>> {
1403 for (region_id, _) in &requests {
1404 self.set_region_status_not_ready(*region_id, &engine, &RegionChange::Catchup);
1405 }
1406 let region_ids = requests
1407 .iter()
1408 .map(|(region_id, _)| *region_id)
1409 .collect::<Vec<_>>();
1410 let mut responses = Vec::with_capacity(requests.len());
1411 match engine
1412 .handle_batch_catchup_requests(parallelism, requests)
1413 .await
1414 {
1415 Ok(results) => {
1416 for (region_id, result) in results {
1417 match result {
1418 Ok(_) => {
1419 if let Err(e) = self
1420 .set_region_status_ready(
1421 region_id,
1422 engine.clone(),
1423 RegionChange::Catchup,
1424 )
1425 .await
1426 {
1427 error!(e; "Failed to set region to ready: {}", region_id);
1428 responses.push((region_id, Err(BoxedError::new(e))));
1429 } else {
1430 responses.push((region_id, Ok(())));
1431 }
1432 }
1433 Err(e) => {
1434 self.unset_region_status(region_id, &engine, RegionChange::Catchup);
1435 error!(e; "Failed to catchup region: {}", region_id);
1436 responses.push((region_id, Err(e)));
1437 }
1438 }
1439 }
1440 }
1441 Err(e) => {
1442 for region_id in region_ids {
1443 self.unset_region_status(region_id, &engine, RegionChange::Catchup);
1444 }
1445 error!(e; "Failed to catchup batch regions");
1446 return error::UnexpectedSnafu {
1447 violated: format!("Failed to catchup batch regions: {:?}", e),
1448 }
1449 .fail();
1450 }
1451 }
1452
1453 Ok(responses)
1454 }
1455
1456 pub async fn handle_batch_catchup_requests(
1457 &self,
1458 parallelism: usize,
1459 requests: Vec<(RegionId, RegionCatchupRequest)>,
1460 ) -> Result<Vec<(RegionId, std::result::Result<(), BoxedError>)>> {
1461 let mut engine_grouped_requests: HashMap<String, Vec<_>> = HashMap::new();
1462
1463 let mut responses = Vec::with_capacity(requests.len());
1464 for (region_id, request) in requests {
1465 if let Ok(engine) = self.get_engine(region_id, &RegionChange::Catchup) {
1466 match engine {
1467 CurrentEngine::Engine(engine) => {
1468 engine_grouped_requests
1469 .entry(engine.name().to_string())
1470 .or_default()
1471 .push((region_id, request));
1472 }
1473 CurrentEngine::EarlyReturn(_) => {
1474 return error::UnexpectedSnafu {
1475 violated: format!("Unexpected engine type for region {}", region_id),
1476 }
1477 .fail();
1478 }
1479 }
1480 } else {
1481 responses.push((
1482 region_id,
1483 Err(BoxedError::new(
1484 error::RegionNotFoundSnafu { region_id }.build(),
1485 )),
1486 ));
1487 }
1488 }
1489
1490 for (engine, requests) in engine_grouped_requests {
1491 let engine = self
1492 .engines
1493 .read()
1494 .unwrap()
1495 .get(&engine)
1496 .with_context(|| RegionEngineNotFoundSnafu { name: &engine })?
1497 .clone();
1498 responses.extend(
1499 self.handle_batch_catchup_requests_inner(engine, parallelism, requests)
1500 .await?,
1501 );
1502 }
1503
1504 Ok(responses)
1505 }
1506
1507 pub async fn handle_batch_request(
1511 &self,
1512 batch_request: BatchRegionDdlRequest,
1513 ) -> Result<RegionResponse> {
1514 let region_changes = match &batch_request {
1515 BatchRegionDdlRequest::Create(requests) => requests
1516 .iter()
1517 .map(|(region_id, create)| {
1518 let attribute = parse_region_attribute(&create.engine, &create.options)?;
1519 Ok((*region_id, RegionChange::Register(attribute)))
1520 })
1521 .collect::<Result<Vec<_>>>()?,
1522 BatchRegionDdlRequest::Drop(requests) => requests
1523 .iter()
1524 .map(|(region_id, _)| (*region_id, RegionChange::Deregisters))
1525 .collect::<Vec<_>>(),
1526 BatchRegionDdlRequest::Alter(requests) => requests
1527 .iter()
1528 .map(|(region_id, _)| (*region_id, RegionChange::None))
1529 .collect::<Vec<_>>(),
1530 };
1531
1532 let (first_region_id, first_region_change) = region_changes.first().unwrap();
1535 let engine = match self.get_engine(*first_region_id, first_region_change)? {
1536 CurrentEngine::Engine(engine) => engine,
1537 CurrentEngine::EarlyReturn(rows) => return Ok(RegionResponse::new(rows)),
1538 };
1539
1540 for (region_id, region_change) in region_changes.iter() {
1541 self.set_region_status_not_ready(*region_id, &engine, region_change);
1542 }
1543
1544 let ddl_type = batch_request.request_type();
1545 let result = engine
1546 .handle_batch_ddl_requests(batch_request)
1547 .await
1548 .context(HandleBatchDdlRequestSnafu { ddl_type });
1549
1550 match result {
1551 Ok(result) => {
1552 for (region_id, region_change) in ®ion_changes {
1553 self.set_region_status_ready(*region_id, engine.clone(), *region_change)
1554 .await?;
1555 }
1556
1557 Ok(RegionResponse {
1558 affected_rows: result.affected_rows,
1559 extensions: result.extensions,
1560 metadata: Vec::new(),
1561 })
1562 }
1563 Err(err) => {
1564 for (region_id, region_change) in region_changes {
1565 self.unset_region_status(region_id, &engine, region_change);
1566 }
1567
1568 Err(err)
1569 }
1570 }
1571 }
1572
1573 pub async fn handle_request(
1574 &self,
1575 region_id: RegionId,
1576 request: RegionRequest,
1577 ) -> Result<RegionResponse> {
1578 let request_type = request.request_type();
1579 let _timer = crate::metrics::HANDLE_REGION_REQUEST_ELAPSED
1580 .with_label_values(&[request_type])
1581 .start_timer();
1582
1583 let region_change = match &request {
1584 RegionRequest::Create(create) => {
1585 let attribute = parse_region_attribute(&create.engine, &create.options)?;
1586 RegionChange::Register(attribute)
1587 }
1588 RegionRequest::Open(open) => {
1589 let attribute = parse_region_attribute(&open.engine, &open.options)?;
1590 RegionChange::Register(attribute)
1591 }
1592 RegionRequest::CleanUp(clean_up) => {
1593 let attribute = parse_region_attribute(&clean_up.engine, &clean_up.options)?;
1594 RegionChange::OfflineCleanup(attribute)
1595 }
1596 RegionRequest::Close(_) | RegionRequest::Drop(_) => RegionChange::Deregisters,
1597 RegionRequest::Put(_) | RegionRequest::Delete(_) | RegionRequest::BulkInserts(_) => {
1598 RegionChange::Ingest
1599 }
1600 RegionRequest::Alter(_)
1601 | RegionRequest::Flush(_)
1602 | RegionRequest::Compact(_)
1603 | RegionRequest::Truncate(_)
1604 | RegionRequest::BuildIndex(_)
1605 | RegionRequest::EnterStaging(_)
1606 | RegionRequest::ApplyStagingManifest(_) => RegionChange::None,
1607 RegionRequest::Catchup(_) => RegionChange::Catchup,
1608 };
1609
1610 let engine = match self.get_engine(region_id, ®ion_change)? {
1611 CurrentEngine::Engine(engine) => engine,
1612 CurrentEngine::EarlyReturn(rows) => return Ok(RegionResponse::new(rows)),
1613 };
1614
1615 self.set_region_status_not_ready(region_id, &engine, ®ion_change);
1617
1618 match engine
1619 .handle_request(region_id, request)
1620 .await
1621 .with_context(|_| HandleRegionRequestSnafu { region_id })
1622 {
1623 Ok(result) => {
1624 if matches!(region_change, RegionChange::Ingest) {
1626 crate::metrics::REGION_CHANGED_ROW_COUNT
1627 .with_label_values(&[request_type])
1628 .inc_by(result.affected_rows as u64);
1629 }
1630 self.set_region_status_ready(region_id, engine.clone(), region_change)
1632 .await?;
1633
1634 Ok(RegionResponse {
1635 affected_rows: result.affected_rows,
1636 extensions: result.extensions,
1637 metadata: Vec::new(),
1638 })
1639 }
1640 Err(err) => {
1641 if matches!(region_change, RegionChange::Ingest) {
1642 crate::metrics::REGION_SERVER_INSERT_FAIL_COUNT
1643 .with_label_values(&[request_type])
1644 .inc();
1645 }
1646 self.unset_region_status(region_id, &engine, region_change);
1648 Err(err)
1649 }
1650 }
1651 }
1652
1653 pub async fn handle_sync_region(
1655 &self,
1656 engine: &RegionEngineRef,
1657 region_id: RegionId,
1658 request: SyncRegionFromRequest,
1659 ) -> Result<()> {
1660 let Some(new_opened_regions) = engine
1661 .sync_region(region_id, request)
1662 .await
1663 .with_context(|_| HandleRegionRequestSnafu { region_id })?
1664 .new_opened_logical_region_ids()
1665 else {
1666 return Ok(());
1667 };
1668
1669 for region in &new_opened_regions {
1670 self.region_map
1671 .insert(*region, RegionEngineWithStatus::Ready(engine.clone()));
1672 }
1673 if !new_opened_regions.is_empty() {
1674 info!(
1675 region_id = %region_id,
1676 logical_region_count = new_opened_regions.len(),
1677 logical_regions = ?new_opened_regions,
1678 "Logical regions are registered"
1679 );
1680 }
1681
1682 Ok(())
1683 }
1684
1685 fn set_region_status_not_ready(
1686 &self,
1687 region_id: RegionId,
1688 engine: &RegionEngineRef,
1689 region_change: &RegionChange,
1690 ) {
1691 match region_change {
1692 RegionChange::Register(_) => {
1693 self.region_map.insert(
1694 region_id,
1695 RegionEngineWithStatus::Registering(engine.clone()),
1696 );
1697 }
1698 RegionChange::Deregisters => {
1699 self.region_map.insert(
1700 region_id,
1701 RegionEngineWithStatus::Deregistering(engine.clone()),
1702 );
1703 }
1704 _ => {}
1705 }
1706 }
1707
1708 fn unset_region_status(
1709 &self,
1710 region_id: RegionId,
1711 engine: &RegionEngineRef,
1712 region_change: RegionChange,
1713 ) {
1714 match region_change {
1715 RegionChange::None | RegionChange::Ingest | RegionChange::OfflineCleanup(_) => {}
1716 RegionChange::Register(_) => {
1717 self.region_map.remove(®ion_id);
1718 }
1719 RegionChange::Deregisters => {
1720 self.region_map
1721 .insert(region_id, RegionEngineWithStatus::Ready(engine.clone()));
1722 }
1723 RegionChange::Catchup => {}
1724 }
1725 }
1726
1727 async fn set_region_status_ready(
1728 &self,
1729 region_id: RegionId,
1730 engine: RegionEngineRef,
1731 region_change: RegionChange,
1732 ) -> Result<()> {
1733 let engine_type = engine.name();
1734 match region_change {
1735 RegionChange::None | RegionChange::Ingest | RegionChange::OfflineCleanup(_) => {}
1736 RegionChange::Register(attribute) => {
1737 info!(
1738 "Region {region_id} is registered to engine {}",
1739 attribute.engine()
1740 );
1741 self.region_map
1742 .insert(region_id, RegionEngineWithStatus::Ready(engine.clone()));
1743
1744 match attribute {
1745 RegionAttribute::Metric { physical } => {
1746 if physical {
1747 self.register_logical_regions(&engine, region_id).await?;
1749 self.event_listener.on_region_registered(region_id);
1751 }
1752 }
1753 RegionAttribute::Mito => self.event_listener.on_region_registered(region_id),
1754 RegionAttribute::File => {
1755 }
1757 }
1758 }
1759 RegionChange::Deregisters => {
1760 info!("Region {region_id} is deregistered from engine {engine_type}");
1761 self.region_map
1762 .remove(®ion_id)
1763 .map(|(id, engine)| engine.set_region_role(id, RegionRole::Follower));
1764 self.event_listener.on_region_deregistered(region_id);
1765 }
1766 RegionChange::Catchup => {
1767 if is_metric_engine(engine.name()) {
1768 self.register_logical_regions(&engine, region_id).await?;
1770 }
1771 }
1772 }
1773 Ok(())
1774 }
1775
1776 async fn register_logical_regions(
1777 &self,
1778 engine: &RegionEngineRef,
1779 physical_region_id: RegionId,
1780 ) -> Result<()> {
1781 let metric_engine =
1782 engine
1783 .as_any()
1784 .downcast_ref::<MetricEngine>()
1785 .context(UnexpectedSnafu {
1786 violated: format!(
1787 "expecting engine type '{}', actual '{}'",
1788 METRIC_ENGINE_NAME,
1789 engine.name(),
1790 ),
1791 })?;
1792
1793 let logical_regions = metric_engine
1794 .logical_regions(physical_region_id)
1795 .await
1796 .context(FindLogicalRegionsSnafu { physical_region_id })?;
1797
1798 for region in &logical_regions {
1799 self.region_map
1800 .insert(*region, RegionEngineWithStatus::Ready(engine.clone()));
1801 }
1802 if !logical_regions.is_empty() {
1803 info!(
1804 physical_region_id = %physical_region_id,
1805 logical_region_count = logical_regions.len(),
1806 logical_regions = ?logical_regions,
1807 "Logical regions are registered"
1808 );
1809 }
1810 Ok(())
1811 }
1812
1813 pub async fn handle_read(
1814 self: &Arc<Self>,
1815 request: QueryRequest,
1816 query_ctx: QueryContextRef,
1817 ) -> Result<SendableRecordBatchStream> {
1818 let live_analyze_metrics =
1819 query_ctx.explain_verbose() && query_ctx.live_analyze_metrics_enabled();
1820 let inner = self.clone();
1821 let mut stream = common_runtime::spawn_query(async move {
1822 inner.handle_read_inner(request, query_ctx).await
1823 })
1824 .await
1825 .context(RuntimeJoinSnafu {
1826 request_type: "read",
1827 })??;
1828 let schema = stream.schema();
1829 let output_ordering = stream.output_ordering().map(|ordering| ordering.to_vec());
1830
1831 let (sender, receiver) = mpsc::channel(QUERY_RUNTIME_STREAM_BUFFER_SIZE);
1832 let metrics = QueryRuntimeStream::metrics_store();
1833 let producer_metrics = metrics.clone();
1834
1835 let producer_handle = common_runtime::spawn_query(async move {
1836 if live_analyze_metrics {
1837 loop {
1838 match time::timeout(FLIGHT_METRICS_HEARTBEAT_INTERVAL, stream.next()).await {
1839 Ok(Some(batch)) => {
1840 *producer_metrics.write().unwrap() = stream.metrics();
1841 if sender.send(batch).await.is_err() {
1842 return;
1843 }
1844 }
1845 Ok(None) => break,
1846 Err(_) => {
1847 *producer_metrics.write().unwrap() = stream.metrics();
1848 }
1849 }
1850 }
1851 } else {
1852 while let Some(batch) = stream.next().await {
1853 *producer_metrics.write().unwrap() = stream.metrics();
1854 if sender.send(batch).await.is_err() {
1855 break;
1856 }
1857 }
1858 }
1859 *producer_metrics.write().unwrap() = stream.metrics();
1860 });
1861
1862 Ok(Box::pin(
1863 QueryRuntimeStream::new(schema, receiver)
1864 .with_output_ordering(output_ordering)
1865 .with_metrics_store(metrics)
1866 .with_producer_handle(producer_handle),
1867 ))
1868 }
1869
1870 async fn handle_read_inner(
1871 &self,
1872 request: QueryRequest,
1873 query_ctx: QueryContextRef,
1874 ) -> Result<SendableRecordBatchStream> {
1875 let result = self
1878 .query_engine
1879 .execute(request.plan, query_ctx)
1880 .await
1881 .context(ExecuteLogicalPlanSnafu)?;
1882
1883 match result.data {
1884 OutputData::AffectedRows(_) | OutputData::RecordBatches(_) => {
1885 UnsupportedOutputSnafu { expected: "stream" }.fail()
1886 }
1887 OutputData::Stream(stream) => Ok(stream),
1888 }
1889 }
1890
1891 async fn stop(&self) -> Result<()> {
1892 let regions = self
1900 .region_map
1901 .iter()
1902 .map(|x| (*x.key(), x.value().clone()))
1903 .collect::<Vec<_>>();
1904 let num_regions = regions.len();
1905
1906 for (region_id, engine) in regions {
1907 let closed = engine
1908 .handle_request(
1909 region_id,
1910 RegionRequest::Close(RegionCloseRequest::default()),
1911 )
1912 .await;
1913 match closed {
1914 Ok(_) => debug!("Region {region_id} is closed"),
1915 Err(e) => warn!("Failed to close region {region_id}, err: {e}"),
1916 }
1917 }
1918 self.region_map.clear();
1919 info!("closed {num_regions} regions");
1920
1921 drop(self.mito_engine.write().unwrap().take());
1922 let engines = self.engines.write().unwrap().drain().collect::<Vec<_>>();
1923 for (engine_name, engine) in engines {
1924 engine
1925 .stop()
1926 .await
1927 .context(StopRegionEngineSnafu { name: &engine_name })?;
1928 info!("Region engine {engine_name} is stopped");
1929 }
1930
1931 Ok(())
1932 }
1933}
1934
1935#[derive(Debug, Clone, Copy)]
1936enum RegionChange {
1937 None,
1938 Register(RegionAttribute),
1939 OfflineCleanup(RegionAttribute),
1940 Deregisters,
1941 Catchup,
1942 Ingest,
1943}
1944
1945fn is_metric_engine(engine: &str) -> bool {
1946 engine == METRIC_ENGINE_NAME
1947}
1948
1949fn parse_region_attribute(
1950 engine: &str,
1951 options: &HashMap<String, String>,
1952) -> Result<RegionAttribute> {
1953 match engine {
1954 MITO_ENGINE_NAME => Ok(RegionAttribute::Mito),
1955 METRIC_ENGINE_NAME => {
1956 let physical = !options.contains_key(LOGICAL_TABLE_METADATA_KEY);
1957
1958 Ok(RegionAttribute::Metric { physical })
1959 }
1960 FILE_ENGINE_NAME => Ok(RegionAttribute::File),
1961 _ => error::UnexpectedSnafu {
1962 violated: format!("Unknown engine: {}", engine),
1963 }
1964 .fail(),
1965 }
1966}
1967
1968#[derive(Debug, Clone, Copy)]
1969enum RegionAttribute {
1970 Mito,
1971 Metric { physical: bool },
1972 File,
1973}
1974
1975impl RegionAttribute {
1976 fn engine(&self) -> &'static str {
1977 match self {
1978 RegionAttribute::Mito => MITO_ENGINE_NAME,
1979 RegionAttribute::Metric { .. } => METRIC_ENGINE_NAME,
1980 RegionAttribute::File => FILE_ENGINE_NAME,
1981 }
1982 }
1983}
1984
1985#[cfg(test)]
1986mod tests {
1987 use std::assert_matches;
1988 use std::collections::HashMap;
1989 use std::sync::Arc;
1990
1991 use api::v1::{Rows, SemanticType};
1992 use common_error::ext::{ErrorExt, RetryHint};
1993 use common_recordbatch::RecordBatches;
1994 use common_recordbatch::adapter::{RecordBatchMetrics, RegionWatermarkEntry};
1995 use datatypes::prelude::{ConcreteDataType, VectorRef};
1996 use datatypes::schema::{ColumnSchema, Schema};
1997 use datatypes::vectors::Int32Vector;
1998 use futures_util::StreamExt;
1999 use mito2::test_util::CreateRequestBuilder;
2000 use query::options::FLOW_RETURN_REGION_SEQ;
2001 use store_api::metadata::{ColumnMetadata, RegionMetadata, RegionMetadataBuilder};
2002 use store_api::region_engine::RegionEngine;
2003 use store_api::region_request::{
2004 PathType, RegionCleanUpRequest, RegionCompactRequest, RegionDeleteRequest,
2005 RegionDropRequest, RegionOpenRequest, RegionPutRequest, RegionTruncateRequest,
2006 };
2007 use store_api::storage::RegionId;
2008
2009 use super::*;
2010 use crate::tests::{MockRegionEngine, mock_region_server};
2011
2012 #[test]
2013 fn test_is_ingest_request() {
2014 let rows = || Rows {
2015 schema: Vec::new(),
2016 rows: Vec::new(),
2017 };
2018
2019 assert!(RegionServerInner::is_ingest_request(&RegionRequest::Put(
2020 RegionPutRequest {
2021 skip_wal: false,
2022 rows: rows(),
2023 hint: None,
2024 partition_expr_version: None,
2025 },
2026 )));
2027 assert!(RegionServerInner::is_ingest_request(
2028 &RegionRequest::Delete(RegionDeleteRequest {
2029 rows: rows(),
2030 hint: None,
2031 partition_expr_version: None,
2032 },)
2033 ));
2034 assert!(!RegionServerInner::is_ingest_request(
2035 &RegionRequest::Compact(RegionCompactRequest::default()),
2036 ));
2037 }
2038
2039 fn single_value_stream() -> SendableRecordBatchStream {
2040 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
2041 "v",
2042 ConcreteDataType::int32_datatype(),
2043 false,
2044 )]));
2045 let values: VectorRef = Arc::new(Int32Vector::from_slice([1]));
2046 let batch = RecordBatch::new(schema.clone(), vec![values]).unwrap();
2047 RecordBatches::try_new(schema, vec![batch])
2048 .unwrap()
2049 .as_stream()
2050 }
2051
2052 #[tokio::test]
2053 async fn test_region_watermark_stream_only_sets_terminal_metrics() {
2054 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
2055 "v",
2056 ConcreteDataType::int32_datatype(),
2057 false,
2058 )]));
2059 let values: VectorRef = Arc::new(Int32Vector::from_slice([1, 2]));
2060 let batch = RecordBatch::new(schema.clone(), vec![values]).unwrap();
2061 let stream = RecordBatches::try_new(schema, vec![batch])
2062 .unwrap()
2063 .as_stream();
2064
2065 let region_id = RegionId::new(42, 7);
2066 let wrapped = RegionWatermarkStream::new(stream, region_id, 99);
2067 let mut pinned = Box::pin(wrapped);
2068
2069 assert!(pinned.as_ref().get_ref().metrics().is_none());
2070 while pinned.next().await.is_some() {}
2071
2072 let metrics = pinned.as_ref().get_ref().metrics().unwrap();
2073 assert_eq!(
2074 metrics.region_watermarks,
2075 vec![RegionWatermarkEntry {
2076 region_id: region_id.as_u64(),
2077 watermark: Some(99),
2078 }]
2079 );
2080 }
2081
2082 #[test]
2083 fn test_region_watermark_stream_preserves_unproved_watermark() {
2084 let schema = Arc::new(Schema::new(vec![ColumnSchema::new(
2085 "v",
2086 ConcreteDataType::int32_datatype(),
2087 false,
2088 )]));
2089 let values: VectorRef = Arc::new(Int32Vector::from_slice([1]));
2090 let batch = RecordBatch::new(schema.clone(), vec![values]).unwrap();
2091 let stream = RecordBatches::try_new(schema, vec![batch])
2092 .unwrap()
2093 .as_stream();
2094
2095 let region_id = RegionId::new(42, 7);
2096 let wrapped = RegionWatermarkStream::new(stream, region_id, 99);
2097 let metrics = RecordBatchMetrics {
2098 region_watermarks: vec![RegionWatermarkEntry {
2099 region_id: region_id.as_u64(),
2100 watermark: None,
2101 }],
2102 ..Default::default()
2103 };
2104
2105 let merged = wrapped.merged_metrics(metrics);
2106 assert_eq!(
2107 merged.region_watermarks,
2108 vec![RegionWatermarkEntry {
2109 region_id: region_id.as_u64(),
2110 watermark: None,
2111 }]
2112 );
2113 }
2114
2115 #[tokio::test]
2116 async fn test_wrap_flow_region_watermark_stream_adds_terminal_metrics() {
2117 let region_id = RegionId::new(42, 7);
2118 let query_ctx = Arc::new(
2119 QueryContextBuilder::default()
2120 .extensions(HashMap::from([(
2121 FLOW_RETURN_REGION_SEQ.to_string(),
2122 "true".to_string(),
2123 )]))
2124 .build(),
2125 );
2126 query_ctx.set_snapshot(region_id.as_u64(), 99);
2127
2128 let wrapped =
2129 wrap_flow_region_watermark_stream(single_value_stream(), region_id, &query_ctx);
2130 let mut pinned = Box::pin(wrapped);
2131 while pinned.next().await.is_some() {}
2132
2133 let metrics = pinned.as_ref().get_ref().metrics().unwrap();
2134 assert_eq!(
2135 metrics.region_watermarks,
2136 vec![RegionWatermarkEntry {
2137 region_id: region_id.as_u64(),
2138 watermark: Some(99),
2139 }]
2140 );
2141 }
2142
2143 #[tokio::test]
2144 async fn test_wrap_flow_region_watermark_stream_skips_without_extension() {
2145 let region_id = RegionId::new(42, 7);
2146 let query_ctx = Arc::new(QueryContextBuilder::default().build());
2147 query_ctx.set_snapshot(region_id.as_u64(), 99);
2148
2149 let wrapped =
2150 wrap_flow_region_watermark_stream(single_value_stream(), region_id, &query_ctx);
2151 let mut pinned = Box::pin(wrapped);
2152 while pinned.next().await.is_some() {}
2153
2154 assert!(pinned.as_ref().get_ref().metrics().is_none());
2155 }
2156
2157 #[tokio::test]
2158 async fn test_wrap_flow_region_watermark_stream_skips_without_snapshot() {
2159 let region_id = RegionId::new(42, 7);
2160 let query_ctx = Arc::new(
2161 QueryContextBuilder::default()
2162 .extensions(HashMap::from([(
2163 FLOW_RETURN_REGION_SEQ.to_string(),
2164 "true".to_string(),
2165 )]))
2166 .build(),
2167 );
2168
2169 let wrapped =
2170 wrap_flow_region_watermark_stream(single_value_stream(), region_id, &query_ctx);
2171 let mut pinned = Box::pin(wrapped);
2172 while pinned.next().await.is_some() {}
2173
2174 assert!(pinned.as_ref().get_ref().metrics().is_none());
2175 }
2176
2177 #[tokio::test]
2178 async fn test_offline_cleanup_does_not_register_region() {
2179 let mut mock_region_server = mock_region_server();
2180 let (engine, mut receiver) = MockRegionEngine::new(MITO_ENGINE_NAME);
2181 mock_region_server.register_engine(engine);
2182
2183 let region_id = RegionId::new(1, 1);
2184 let response = mock_region_server
2185 .handle_request(
2186 region_id,
2187 RegionRequest::CleanUp(RegionCleanUpRequest {
2188 engine: MITO_ENGINE_NAME.to_string(),
2189 table_dir: String::new(),
2190 path_type: PathType::Bare,
2191 options: HashMap::new(),
2192 }),
2193 )
2194 .await
2195 .unwrap();
2196
2197 assert_eq!(response.affected_rows, 0);
2198 let (handled_region_id, handled_request) = receiver.try_recv().unwrap();
2199 assert_eq!(handled_region_id, region_id);
2200 assert_matches!(handled_request, RegionRequest::CleanUp(_));
2201 assert!(
2202 mock_region_server
2203 .inner
2204 .region_map
2205 .get(®ion_id)
2206 .is_none()
2207 );
2208 }
2209
2210 #[tokio::test]
2211 async fn test_offline_cleanup_rejects_registered_region() {
2212 let mut mock_region_server = mock_region_server();
2213 let (engine, mut receiver) = MockRegionEngine::new(MITO_ENGINE_NAME);
2214 mock_region_server.register_engine(engine.clone());
2215
2216 let region_id = RegionId::new(1, 1);
2217 mock_region_server
2218 .inner
2219 .region_map
2220 .insert(region_id, RegionEngineWithStatus::Ready(engine));
2221
2222 let err = mock_region_server
2223 .handle_request(
2224 region_id,
2225 RegionRequest::CleanUp(RegionCleanUpRequest {
2226 engine: MITO_ENGINE_NAME.to_string(),
2227 table_dir: String::new(),
2228 path_type: PathType::Bare,
2229 options: HashMap::new(),
2230 }),
2231 )
2232 .await
2233 .unwrap_err();
2234
2235 assert_eq!(err.status_code(), StatusCode::RegionBusy);
2236 assert!(receiver.try_recv().is_err());
2237 assert!(matches!(
2238 mock_region_server
2239 .inner
2240 .region_map
2241 .get(®ion_id)
2242 .unwrap()
2243 .clone(),
2244 RegionEngineWithStatus::Ready(_)
2245 ));
2246 }
2247
2248 #[tokio::test]
2249 async fn test_region_registering() {
2250 common_telemetry::init_default_ut_logging();
2251
2252 let mut mock_region_server = mock_region_server();
2253 let (engine, _receiver) = MockRegionEngine::new(MITO_ENGINE_NAME);
2254 let engine_name = engine.name();
2255 mock_region_server.register_engine(engine.clone());
2256 let region_id = RegionId::new(1, 1);
2257 let builder = CreateRequestBuilder::new();
2258 let create_req = builder.build();
2259 mock_region_server.inner.region_map.insert(
2261 region_id,
2262 RegionEngineWithStatus::Registering(engine.clone()),
2263 );
2264 let response = mock_region_server
2265 .handle_request(region_id, RegionRequest::Create(create_req))
2266 .await
2267 .unwrap();
2268 assert_eq!(response.affected_rows, 0);
2269 let status = mock_region_server
2270 .inner
2271 .region_map
2272 .get(®ion_id)
2273 .unwrap()
2274 .clone();
2275 assert!(matches!(status, RegionEngineWithStatus::Ready(_)));
2276
2277 mock_region_server.inner.region_map.insert(
2278 region_id,
2279 RegionEngineWithStatus::Registering(engine.clone()),
2280 );
2281 let response = mock_region_server
2282 .handle_request(
2283 region_id,
2284 RegionRequest::Open(RegionOpenRequest {
2285 engine: engine_name.to_string(),
2286 table_dir: String::new(),
2287 path_type: PathType::Bare,
2288 options: Default::default(),
2289 skip_wal_replay: false,
2290 checkpoint: None,
2291 requirements: Default::default(),
2292 }),
2293 )
2294 .await
2295 .unwrap();
2296 assert_eq!(response.affected_rows, 0);
2297 let status = mock_region_server
2298 .inner
2299 .region_map
2300 .get(®ion_id)
2301 .unwrap()
2302 .clone();
2303 assert!(matches!(status, RegionEngineWithStatus::Ready(_)));
2304 }
2305
2306 #[tokio::test]
2307 async fn test_region_deregistering() {
2308 common_telemetry::init_default_ut_logging();
2309
2310 let mut mock_region_server = mock_region_server();
2311 let (engine, _receiver) = MockRegionEngine::new(MITO_ENGINE_NAME);
2312
2313 mock_region_server.register_engine(engine.clone());
2314
2315 let region_id = RegionId::new(1, 1);
2316
2317 mock_region_server.inner.region_map.insert(
2319 region_id,
2320 RegionEngineWithStatus::Deregistering(engine.clone()),
2321 );
2322
2323 let response = mock_region_server
2324 .handle_request(
2325 region_id,
2326 RegionRequest::Drop(RegionDropRequest {
2327 fast_path: false,
2328 force: false,
2329 partial_drop: false,
2330 }),
2331 )
2332 .await
2333 .unwrap();
2334 assert_eq!(response.affected_rows, 0);
2335
2336 let status = mock_region_server
2337 .inner
2338 .region_map
2339 .get(®ion_id)
2340 .unwrap()
2341 .clone();
2342 assert!(matches!(status, RegionEngineWithStatus::Deregistering(_)));
2343
2344 mock_region_server.inner.region_map.insert(
2345 region_id,
2346 RegionEngineWithStatus::Deregistering(engine.clone()),
2347 );
2348
2349 let response = mock_region_server
2350 .handle_request(
2351 region_id,
2352 RegionRequest::Close(RegionCloseRequest::default()),
2353 )
2354 .await
2355 .unwrap();
2356 assert_eq!(response.affected_rows, 0);
2357
2358 let status = mock_region_server
2359 .inner
2360 .region_map
2361 .get(®ion_id)
2362 .unwrap()
2363 .clone();
2364 assert!(matches!(status, RegionEngineWithStatus::Deregistering(_)));
2365 }
2366
2367 #[tokio::test]
2368 async fn test_region_not_ready() {
2369 common_telemetry::init_default_ut_logging();
2370
2371 let mut mock_region_server = mock_region_server();
2372 let (engine, _receiver) = MockRegionEngine::new(MITO_ENGINE_NAME);
2373
2374 mock_region_server.register_engine(engine.clone());
2375
2376 let region_id = RegionId::new(1, 1);
2377
2378 mock_region_server.inner.region_map.insert(
2380 region_id,
2381 RegionEngineWithStatus::Registering(engine.clone()),
2382 );
2383
2384 let err = mock_region_server
2385 .handle_request(
2386 region_id,
2387 RegionRequest::Truncate(RegionTruncateRequest::All),
2388 )
2389 .await
2390 .unwrap_err();
2391
2392 assert_eq!(err.status_code(), StatusCode::RegionNotReady);
2393 }
2394
2395 #[tokio::test]
2396 async fn test_region_request_failed() {
2397 common_telemetry::init_default_ut_logging();
2398
2399 let mut mock_region_server = mock_region_server();
2400 let (engine, _receiver) = MockRegionEngine::with_mock_fn(
2401 MITO_ENGINE_NAME,
2402 Box::new(|_region_id, _request| {
2403 error::UnexpectedSnafu {
2404 violated: "test".to_string(),
2405 }
2406 .fail()
2407 }),
2408 );
2409
2410 mock_region_server.register_engine(engine.clone());
2411
2412 let region_id = RegionId::new(1, 1);
2413 let builder = CreateRequestBuilder::new();
2414 let create_req = builder.build();
2415 mock_region_server
2416 .handle_request(region_id, RegionRequest::Create(create_req))
2417 .await
2418 .unwrap_err();
2419
2420 let status = mock_region_server.inner.region_map.get(®ion_id);
2421 assert!(status.is_none());
2422
2423 mock_region_server
2424 .inner
2425 .region_map
2426 .insert(region_id, RegionEngineWithStatus::Ready(engine.clone()));
2427
2428 mock_region_server
2429 .handle_request(
2430 region_id,
2431 RegionRequest::Drop(RegionDropRequest {
2432 fast_path: false,
2433 force: false,
2434 partial_drop: false,
2435 }),
2436 )
2437 .await
2438 .unwrap_err();
2439
2440 let status = mock_region_server.inner.region_map.get(®ion_id);
2441 assert!(status.is_some());
2442 }
2443
2444 #[tokio::test]
2445 async fn test_batch_open_region_ignore_nonexistent_regions() {
2446 common_telemetry::init_default_ut_logging();
2447 let mut mock_region_server = mock_region_server();
2448 let (engine, _receiver) = MockRegionEngine::with_mock_fn(
2449 MITO_ENGINE_NAME,
2450 Box::new(|region_id, _request| {
2451 if region_id == RegionId::new(1, 1) {
2452 error::RegionNotFoundSnafu { region_id }.fail()
2453 } else {
2454 Ok(0)
2455 }
2456 }),
2457 );
2458 mock_region_server.register_engine(engine.clone());
2459
2460 let region_ids = mock_region_server
2461 .handle_batch_open_requests(
2462 8,
2463 vec![
2464 (
2465 RegionId::new(1, 1),
2466 RegionOpenRequest {
2467 engine: MITO_ENGINE_NAME.to_string(),
2468 table_dir: String::new(),
2469 path_type: PathType::Bare,
2470 options: Default::default(),
2471 skip_wal_replay: false,
2472 checkpoint: None,
2473 requirements: Default::default(),
2474 },
2475 ),
2476 (
2477 RegionId::new(1, 2),
2478 RegionOpenRequest {
2479 engine: MITO_ENGINE_NAME.to_string(),
2480 table_dir: String::new(),
2481 path_type: PathType::Bare,
2482 options: Default::default(),
2483 skip_wal_replay: false,
2484 checkpoint: None,
2485 requirements: Default::default(),
2486 },
2487 ),
2488 ],
2489 true,
2490 )
2491 .await
2492 .unwrap();
2493 assert_eq!(region_ids, vec![RegionId::new(1, 2)]);
2494
2495 let err = mock_region_server
2496 .handle_batch_open_requests(
2497 8,
2498 vec![
2499 (
2500 RegionId::new(1, 1),
2501 RegionOpenRequest {
2502 engine: MITO_ENGINE_NAME.to_string(),
2503 table_dir: String::new(),
2504 path_type: PathType::Bare,
2505 options: Default::default(),
2506 skip_wal_replay: false,
2507 checkpoint: None,
2508 requirements: Default::default(),
2509 },
2510 ),
2511 (
2512 RegionId::new(1, 2),
2513 RegionOpenRequest {
2514 engine: MITO_ENGINE_NAME.to_string(),
2515 table_dir: String::new(),
2516 path_type: PathType::Bare,
2517 options: Default::default(),
2518 skip_wal_replay: false,
2519 checkpoint: None,
2520 requirements: Default::default(),
2521 },
2522 ),
2523 ],
2524 false,
2525 )
2526 .await
2527 .unwrap_err();
2528 assert_matches!(&err, error::Error::HandleBatchOpenRequest { .. });
2529 assert_eq!(err.status_code(), StatusCode::RegionNotFound);
2530 assert_eq!(err.retry_hint(), RetryHint::NonRetryable);
2531 }
2532
2533 struct CurrentEngineTest {
2534 region_id: RegionId,
2535 current_region_status: Option<RegionEngineWithStatus>,
2536 region_change: RegionChange,
2537 assert: Box<dyn FnOnce(Result<CurrentEngine>)>,
2538 }
2539
2540 #[tokio::test]
2541 async fn test_current_engine() {
2542 common_telemetry::init_default_ut_logging();
2543
2544 let mut mock_region_server = mock_region_server();
2545 let (engine, _) = MockRegionEngine::new(MITO_ENGINE_NAME);
2546 mock_region_server.register_engine(engine.clone());
2547
2548 let region_id = RegionId::new(1024, 1);
2549 let tests = vec![
2550 CurrentEngineTest {
2552 region_id,
2553 current_region_status: None,
2554 region_change: RegionChange::None,
2555 assert: Box::new(|result| {
2556 let err = result.unwrap_err();
2557 assert_eq!(err.status_code(), StatusCode::RegionNotFound);
2558 }),
2559 },
2560 CurrentEngineTest {
2561 region_id,
2562 current_region_status: Some(RegionEngineWithStatus::Ready(engine.clone())),
2563 region_change: RegionChange::None,
2564 assert: Box::new(|result| {
2565 let current_engine = result.unwrap();
2566 assert_matches!(current_engine, CurrentEngine::Engine(_));
2567 }),
2568 },
2569 CurrentEngineTest {
2570 region_id,
2571 current_region_status: Some(RegionEngineWithStatus::Registering(engine.clone())),
2572 region_change: RegionChange::None,
2573 assert: Box::new(|result| {
2574 let err = result.unwrap_err();
2575 assert_eq!(err.status_code(), StatusCode::RegionNotReady);
2576 }),
2577 },
2578 CurrentEngineTest {
2579 region_id,
2580 current_region_status: Some(RegionEngineWithStatus::Deregistering(engine.clone())),
2581 region_change: RegionChange::None,
2582 assert: Box::new(|result| {
2583 let err = result.unwrap_err();
2584 assert_eq!(err.status_code(), StatusCode::RegionNotFound);
2585 }),
2586 },
2587 CurrentEngineTest {
2589 region_id,
2590 current_region_status: None,
2591 region_change: RegionChange::Register(RegionAttribute::Mito),
2592 assert: Box::new(|result| {
2593 let current_engine = result.unwrap();
2594 assert_matches!(current_engine, CurrentEngine::Engine(_));
2595 }),
2596 },
2597 CurrentEngineTest {
2598 region_id,
2599 current_region_status: Some(RegionEngineWithStatus::Registering(engine.clone())),
2600 region_change: RegionChange::Register(RegionAttribute::Mito),
2601 assert: Box::new(|result| {
2602 let current_engine = result.unwrap();
2603 assert_matches!(current_engine, CurrentEngine::Engine(_));
2604 }),
2605 },
2606 CurrentEngineTest {
2607 region_id,
2608 current_region_status: Some(RegionEngineWithStatus::Deregistering(engine.clone())),
2609 region_change: RegionChange::Register(RegionAttribute::Mito),
2610 assert: Box::new(|result| {
2611 let err = result.unwrap_err();
2612 assert_eq!(err.status_code(), StatusCode::RegionBusy);
2613 }),
2614 },
2615 CurrentEngineTest {
2616 region_id,
2617 current_region_status: Some(RegionEngineWithStatus::Ready(engine.clone())),
2618 region_change: RegionChange::Register(RegionAttribute::Mito),
2619 assert: Box::new(|result| {
2620 let current_engine = result.unwrap();
2621 assert_matches!(current_engine, CurrentEngine::Engine(_));
2622 }),
2623 },
2624 CurrentEngineTest {
2626 region_id,
2627 current_region_status: None,
2628 region_change: RegionChange::Deregisters,
2629 assert: Box::new(|result| {
2630 let current_engine = result.unwrap();
2631 assert_matches!(current_engine, CurrentEngine::EarlyReturn(_));
2632 }),
2633 },
2634 CurrentEngineTest {
2635 region_id,
2636 current_region_status: Some(RegionEngineWithStatus::Registering(engine.clone())),
2637 region_change: RegionChange::Deregisters,
2638 assert: Box::new(|result| {
2639 let err = result.unwrap_err();
2640 assert_eq!(err.status_code(), StatusCode::RegionBusy);
2641 }),
2642 },
2643 CurrentEngineTest {
2644 region_id,
2645 current_region_status: Some(RegionEngineWithStatus::Deregistering(engine.clone())),
2646 region_change: RegionChange::Deregisters,
2647 assert: Box::new(|result| {
2648 let current_engine = result.unwrap();
2649 assert_matches!(current_engine, CurrentEngine::EarlyReturn(_));
2650 }),
2651 },
2652 CurrentEngineTest {
2653 region_id,
2654 current_region_status: Some(RegionEngineWithStatus::Ready(engine.clone())),
2655 region_change: RegionChange::Deregisters,
2656 assert: Box::new(|result| {
2657 let current_engine = result.unwrap();
2658 assert_matches!(current_engine, CurrentEngine::Engine(_));
2659 }),
2660 },
2661 CurrentEngineTest {
2663 region_id,
2664 current_region_status: None,
2665 region_change: RegionChange::OfflineCleanup(RegionAttribute::Mito),
2666 assert: Box::new(|result| {
2667 let current_engine = result.unwrap();
2668 assert_matches!(current_engine, CurrentEngine::Engine(_));
2669 }),
2670 },
2671 CurrentEngineTest {
2672 region_id,
2673 current_region_status: Some(RegionEngineWithStatus::Registering(engine.clone())),
2674 region_change: RegionChange::OfflineCleanup(RegionAttribute::Mito),
2675 assert: Box::new(|result| {
2676 let err = result.unwrap_err();
2677 assert_eq!(err.status_code(), StatusCode::RegionBusy);
2678 }),
2679 },
2680 CurrentEngineTest {
2681 region_id,
2682 current_region_status: Some(RegionEngineWithStatus::Deregistering(engine.clone())),
2683 region_change: RegionChange::OfflineCleanup(RegionAttribute::Mito),
2684 assert: Box::new(|result| {
2685 let err = result.unwrap_err();
2686 assert_eq!(err.status_code(), StatusCode::RegionBusy);
2687 }),
2688 },
2689 CurrentEngineTest {
2690 region_id,
2691 current_region_status: Some(RegionEngineWithStatus::Ready(engine.clone())),
2692 region_change: RegionChange::OfflineCleanup(RegionAttribute::Mito),
2693 assert: Box::new(|result| {
2694 let err = result.unwrap_err();
2695 assert_eq!(err.status_code(), StatusCode::RegionBusy);
2696 }),
2697 },
2698 ];
2699
2700 for test in tests {
2701 let CurrentEngineTest {
2702 region_id,
2703 current_region_status,
2704 region_change,
2705 assert,
2706 } = test;
2707
2708 if let Some(status) = current_region_status {
2710 mock_region_server
2711 .inner
2712 .region_map
2713 .insert(region_id, status);
2714 } else {
2715 mock_region_server.inner.region_map.remove(®ion_id);
2716 }
2717
2718 let result = mock_region_server
2719 .inner
2720 .get_engine(region_id, ®ion_change);
2721
2722 assert(result);
2723 }
2724 }
2725
2726 #[tokio::test]
2727 async fn test_region_server_parallelism() {
2728 let p = RegionServerParallelism::from_opts(2, Duration::from_millis(1)).unwrap();
2729 let first_query = p.acquire().await;
2730 assert!(first_query.is_ok());
2731 let second_query = p.acquire().await;
2732 assert!(second_query.is_ok());
2733 let third_query = p.acquire().await;
2734 assert!(third_query.is_err());
2735 let err = third_query.unwrap_err();
2736 assert_eq!(
2737 err.output_msg(),
2738 "Failed to acquire permit under timeouts: deadline has elapsed".to_string()
2739 );
2740 drop(first_query);
2741 let forth_query = p.acquire().await;
2742 assert!(forth_query.is_ok());
2743 }
2744
2745 fn mock_region_metadata(region_id: RegionId) -> RegionMetadata {
2746 let mut metadata_builder = RegionMetadataBuilder::new(region_id);
2747 metadata_builder.push_column_metadata(ColumnMetadata {
2748 column_schema: datatypes::schema::ColumnSchema::new(
2749 "timestamp",
2750 ConcreteDataType::timestamp_nanosecond_datatype(),
2751 false,
2752 ),
2753 semantic_type: SemanticType::Timestamp,
2754 column_id: 0,
2755 });
2756 metadata_builder.push_column_metadata(ColumnMetadata {
2757 column_schema: datatypes::schema::ColumnSchema::new(
2758 "file",
2759 ConcreteDataType::string_datatype(),
2760 true,
2761 ),
2762 semantic_type: SemanticType::Tag,
2763 column_id: 1,
2764 });
2765 metadata_builder.push_column_metadata(ColumnMetadata {
2766 column_schema: datatypes::schema::ColumnSchema::new(
2767 "message",
2768 ConcreteDataType::string_datatype(),
2769 true,
2770 ),
2771 semantic_type: SemanticType::Field,
2772 column_id: 2,
2773 });
2774 metadata_builder.primary_key(vec![1]);
2775 metadata_builder.build().unwrap()
2776 }
2777
2778 #[tokio::test]
2779 async fn test_handle_list_metadata_request() {
2780 common_telemetry::init_default_ut_logging();
2781
2782 let mut mock_region_server = mock_region_server();
2783 let region_id_1 = RegionId::new(1, 0);
2784 let region_id_2 = RegionId::new(2, 0);
2785
2786 let metadata_1 = mock_region_metadata(region_id_1);
2787 let metadata_2 = mock_region_metadata(region_id_2);
2788 let metadatas = vec![Some(metadata_1.clone()), Some(metadata_2.clone())];
2789
2790 let metadata_1 = Arc::new(metadata_1);
2791 let metadata_2 = Arc::new(metadata_2);
2792 let (engine, _) = MockRegionEngine::with_metadata_mock_fn(
2793 MITO_ENGINE_NAME,
2794 Box::new(move |region_id| {
2795 if region_id == region_id_1 {
2796 Ok(metadata_1.clone())
2797 } else if region_id == region_id_2 {
2798 Ok(metadata_2.clone())
2799 } else {
2800 error::RegionNotFoundSnafu { region_id }.fail()
2801 }
2802 }),
2803 );
2804
2805 mock_region_server.register_engine(engine.clone());
2806 mock_region_server
2807 .inner
2808 .region_map
2809 .insert(region_id_1, RegionEngineWithStatus::Ready(engine.clone()));
2810 mock_region_server
2811 .inner
2812 .region_map
2813 .insert(region_id_2, RegionEngineWithStatus::Ready(engine.clone()));
2814
2815 let list_metadata_request = ListMetadataRequest {
2817 region_ids: vec![region_id_1.as_u64(), region_id_2.as_u64()],
2818 };
2819 let response = mock_region_server
2820 .handle_list_metadata_request(&list_metadata_request)
2821 .await
2822 .unwrap();
2823 let decoded_metadata: Vec<Option<RegionMetadata>> =
2824 serde_json::from_slice(&response.metadata).unwrap();
2825 assert_eq!(metadatas, decoded_metadata);
2826 }
2827
2828 #[tokio::test]
2829 async fn test_handle_list_metadata_not_found() {
2830 common_telemetry::init_default_ut_logging();
2831
2832 let mut mock_region_server = mock_region_server();
2833 let region_id_1 = RegionId::new(1, 0);
2834 let region_id_2 = RegionId::new(2, 0);
2835
2836 let metadata_1 = mock_region_metadata(region_id_1);
2837 let metadatas = vec![Some(metadata_1.clone()), None];
2838
2839 let metadata_1 = Arc::new(metadata_1);
2840 let (engine, _) = MockRegionEngine::with_metadata_mock_fn(
2841 MITO_ENGINE_NAME,
2842 Box::new(move |region_id| {
2843 if region_id == region_id_1 {
2844 Ok(metadata_1.clone())
2845 } else {
2846 error::RegionNotFoundSnafu { region_id }.fail()
2847 }
2848 }),
2849 );
2850
2851 mock_region_server.register_engine(engine.clone());
2852 mock_region_server
2853 .inner
2854 .region_map
2855 .insert(region_id_1, RegionEngineWithStatus::Ready(engine.clone()));
2856
2857 let list_metadata_request = ListMetadataRequest {
2859 region_ids: vec![region_id_1.as_u64(), region_id_2.as_u64()],
2860 };
2861 let response = mock_region_server
2862 .handle_list_metadata_request(&list_metadata_request)
2863 .await
2864 .unwrap();
2865 let decoded_metadata: Vec<Option<RegionMetadata>> =
2866 serde_json::from_slice(&response.metadata).unwrap();
2867 assert_eq!(metadatas, decoded_metadata);
2868
2869 mock_region_server
2871 .inner
2872 .region_map
2873 .insert(region_id_2, RegionEngineWithStatus::Ready(engine.clone()));
2874 let response = mock_region_server
2875 .handle_list_metadata_request(&list_metadata_request)
2876 .await
2877 .unwrap();
2878 let decoded_metadata: Vec<Option<RegionMetadata>> =
2879 serde_json::from_slice(&response.metadata).unwrap();
2880 assert_eq!(metadatas, decoded_metadata);
2881 }
2882
2883 #[tokio::test]
2884 async fn test_handle_list_metadata_failed() {
2885 common_telemetry::init_default_ut_logging();
2886
2887 let mut mock_region_server = mock_region_server();
2888 let region_id_1 = RegionId::new(1, 0);
2889
2890 let (engine, _) = MockRegionEngine::with_metadata_mock_fn(
2891 MITO_ENGINE_NAME,
2892 Box::new(move |region_id| {
2893 error::UnexpectedSnafu {
2894 violated: format!("Failed to get region {region_id}"),
2895 }
2896 .fail()
2897 }),
2898 );
2899
2900 mock_region_server.register_engine(engine.clone());
2901 mock_region_server
2902 .inner
2903 .region_map
2904 .insert(region_id_1, RegionEngineWithStatus::Ready(engine.clone()));
2905
2906 let list_metadata_request = ListMetadataRequest {
2908 region_ids: vec![region_id_1.as_u64()],
2909 };
2910 mock_region_server
2911 .handle_list_metadata_request(&list_metadata_request)
2912 .await
2913 .unwrap_err();
2914 }
2915}