Skip to main content

datanode/
region_server.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15mod 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    /// Registers an engine.
169    pub fn register_engine(&mut self, engine: RegionEngineRef) {
170        self.inner.register_engine(engine);
171    }
172
173    /// Sets the topic stats.
174    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    /// Finds the region's engine by its id. If the region is not ready, returns `None`.
179    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    /// Gets the MitoEngine if it's registered.
189    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    /// Returns a table provider for the region. Will set snapshot sequence if available in the context.
252    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(&region_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    /// Handle reads from remote. They're often query requests received by our Arrow Flight service.
279    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    /// Returns all opened and reportable regions.
389    ///
390    /// Notes: except all metrics regions.
391    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                // Filters out any regions whose role equals None.
398                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    /// Returns the reportable topics.
408    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(&region_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(&region_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    /// Set region role state gracefully.
448    ///
449    /// For [SettableRegionRoleState::Follower]:
450    /// After the call returns, the engine ensures that
451    /// no **further** write or flush operations will succeed in this region.
452    ///
453    /// For [SettableRegionRoleState::DowngradingLeader]:
454    /// After the call returns, the engine ensures that
455    /// no **further** write operations will succeed in this region.
456    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(&region_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(&region_id) {
476            Some(e) => e.region_statistic(region_id),
477            None => None,
478        }
479    }
480
481    /// Stop the region server.
482    pub async fn stop(&self) -> Result<()> {
483        self.inner.stop().await
484    }
485
486    #[cfg(test)]
487    /// Registers a region for test purpose.
488    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        // Safety: we have already checked the request type in `RegionServer::handle()`.
507        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        // Try to optimize batch Put requests for metric engine
527        // Returns either Some(response) or None(requests_back)
528        match self.try_handle_metric_batch_puts(requests).await? {
529            Either::Left(response) => Ok(response),
530            Either::Right(requests) => {
531                // Fallback: original parallel processing
532                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    /// Attempts to optimize batch Put requests for metric engine.
596    ///
597    /// Returns Either::Left(response) if optimization succeeded,
598    /// or Either::Right(original_requests) to fall back to parallel processing.
599    ///
600    /// This avoids cloning requests when optimization cannot be applied.
601    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        // Quick check: verify first request is Put and is metric engine
610        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        // SAFETY: If the first request belongs to metric engine, then ALL requests
617        // in this batch are guaranteed to belong to metric engine. This invariant
618        // is maintained by the request batching logic upstream.
619        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        // Check if ALL requests are Put (now we know it's worth checking)
632        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        // Now extract Put requests by consuming ownership (zero clone!)
645        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        // Downcast to MetricEngine and call batch API
654        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    /// Handles the ListMetadata request and retrieves metadata for specified regions.
730    ///
731    /// Returns the results as a JSON-serialized list in the [RegionResponse]. It serializes
732    /// non-existing regions as `null`.
733    #[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        // Collect metadata for each region
740        for region_id in &request.region_ids {
741            let region_id = RegionId::from_u64(*region_id);
742            // Get the engine.
743            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        // Serialize metadata to JSON
764        let json_result = serde_json::to_vec(&region_metadatas).context(SerializeJsonSnafu)?;
765
766        let response = RegionResponse::from_metadata(json_result);
767
768        Ok(response)
769    }
770
771    /// Sync region manifest and registers new opened logical regions.
772    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    /// Remaps manifests from old regions to new regions.
793    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
837/// Wraps a region read stream so terminal metrics can carry the scan-open watermark.
838struct 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    // An opening, or creating region.
1010    Registering(RegionEngineRef),
1011    // A closing, or dropping region.
1012    Deregistering(RegionEngineRef),
1013    // A ready region.
1014    Ready(RegionEngineRef),
1015}
1016
1017impl RegionEngineWithStatus {
1018    /// Returns [RegionEngineRef].
1019    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    /// The number of queries allowed to be executed at the same time.
1048    /// Act as last line of defense on datanode to prevent query overloading.
1049    parallelism: Option<RegionServerParallelism>,
1050    /// The topic stats reporter.
1051    topic_stats_reporter: RwLock<Option<Box<dyn TopicStatsReporter>>>,
1052    /// HACK(zhongzc): Direct MitoEngine handle for diagnostics. This couples the
1053    /// server with a concrete engine; acceptable for now to fetch Mito-specific
1054    /// info (e.g., list SSTs). Consider a diagnostics trait later.
1055    mito_engine: RwLock<Option<MitoEngine>>,
1056    /// TODO(remote-dyn-filter): Reap this query-scoped placeholder registry on query finish/cancel
1057    /// and later fold it into the real remote dyn filter runtime state lifecycle.
1058    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
1088/// Wraps a record batch stream and holds a concurrency permit until the stream is
1089/// fully consumed (dropped), so `max_concurrent_queries` bounds the number of
1090/// in-flight read streams, not just query planning.
1091struct 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
1122/// Wraps `stream` so it holds `permit` until fully consumed. Returns `stream`
1123/// unchanged when no permit was acquired (limiter disabled).
1124fn 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(&region_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 (&region_id, region_change) in &region_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 = &region_changes[&region_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 (&region_id, region_change) in &region_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            // Preserve the first region error so callers can honor its status code and retry hint.
1347            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    // Handle requests in batch.
1508    //
1509    // limitation: all create requests must be in the same engine.
1510    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        // The ddl procedure will ensure all requests are in the same engine.
1533        // Therefore, we can get the engine from the first request.
1534        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 &region_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, &region_change)? {
1611            CurrentEngine::Engine(engine) => engine,
1612            CurrentEngine::EarlyReturn(rows) => return Ok(RegionResponse::new(rows)),
1613        };
1614
1615        // Sets corresponding region status to registering/deregistering before the operation.
1616        self.set_region_status_not_ready(region_id, &engine, &region_change);
1617
1618        match engine
1619            .handle_request(region_id, request)
1620            .await
1621            .with_context(|_| HandleRegionRequestSnafu { region_id })
1622        {
1623            Ok(result) => {
1624                // Update metrics
1625                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                // Sets corresponding region status to ready.
1631                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                // Removes the region status if the operation fails.
1647                self.unset_region_status(region_id, &engine, region_change);
1648                Err(err)
1649            }
1650        }
1651    }
1652
1653    /// Handles the sync region request.
1654    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(&region_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                            // Registers the logical regions belong to the physical region (`region_id`).
1748                            self.register_logical_regions(&engine, region_id).await?;
1749                            // We only send the `on_region_registered` event of the physical region.
1750                            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                        // do nothing
1756                    }
1757                }
1758            }
1759            RegionChange::Deregisters => {
1760                info!("Region {region_id} is deregistered from engine {engine_type}");
1761                self.region_map
1762                    .remove(&region_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                    // Registers the logical regions belong to the physical region (`region_id`).
1769                    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        // TODO(ruihang): add metrics and set trace id
1876
1877        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        // Calling async functions while iterating inside the Dashmap could easily cause the Rust
1893        // complains "higher-ranked lifetime error". Rust can't prove some future is legit.
1894        // Possible related issue: https://github.com/rust-lang/rust/issues/102211
1895        //
1896        // The workaround is to put the async functions in the `common_runtime::spawn_global`. Or like
1897        // it here, collect the values first then use later separately.
1898
1899        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(&region_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(&region_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        // Tries to create/open a registering region.
2260        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(&region_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(&region_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        // Tries to drop/close a registering region.
2318        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(&region_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(&region_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        // Tries to drop/close a registering region.
2379        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(&region_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(&region_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            // RegionChange::None
2551            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            // RegionChange::Register
2588            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            // RegionChange::Deregister
2625            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            // RegionChange::OfflineCleanup
2662            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            // Sets up
2709            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(&region_id);
2716            }
2717
2718            let result = mock_region_server
2719                .inner
2720                .get_engine(region_id, &region_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        // All regions exist.
2816        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        // Not in region map.
2858        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        // Not in region engine.
2870        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        // Failed to get.
2907        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}