Skip to main content

meta_srv/metasrv/
builder.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
15use std::path::Path;
16use std::sync::atomic::AtomicBool;
17use std::sync::{Arc, Mutex, RwLock};
18use std::time::Duration;
19
20use client::client_manager::NodeClients;
21use client::inserter::InsertOptions;
22use common_base::Plugins;
23use common_catalog::consts::{MIN_USER_FLOW_ID, MIN_USER_TABLE_ID};
24use common_event_recorder::{DEFAULT_COMPACTION_TIME_WINDOW, EventRecorderImpl, EventRecorderRef};
25use common_meta::ddl::flow_meta::FlowMetadataAllocator;
26use common_meta::ddl::table_meta::{TableMetadataAllocator, TableMetadataAllocatorRef};
27use common_meta::ddl::{
28    DdlContext, NoopRegionFailureDetectorControl, RegionFailureDetectorControllerRef,
29};
30use common_meta::ddl_manager::{
31    DdlManager, DdlManagerConfiguratorRef, RepartitionProcedureFactoryRef,
32};
33use common_meta::distributed_time_constants::default_distributed_time_constants;
34use common_meta::key::TableMetadataManager;
35use common_meta::key::flow::FlowMetadataManager;
36use common_meta::key::flow::flow_state::FlowStateManager;
37use common_meta::key::runtime_switch::{RuntimeSwitchManager, RuntimeSwitchManagerRef};
38use common_meta::kv_backend::memory::MemoryKvBackend;
39use common_meta::kv_backend::{KvBackendRef, ResettableKvBackendRef};
40use common_meta::node_manager::NodeManagerRef;
41use common_meta::reconciliation::manager::ReconciliationManager;
42use common_meta::region_keeper::MemoryRegionKeeper;
43use common_meta::region_registry::LeaderRegionRegistry;
44use common_meta::sequence::SequenceBuilder;
45use common_meta::state_store::KvStateStore;
46use common_meta::stats::topic::TopicStatsRegistry;
47use common_meta::wal_provider::{build_kafka_client, build_wal_provider};
48use common_procedure::ProcedureManagerRef;
49use common_procedure::local::{LocalManager, ManagerConfig};
50use common_stat::ResourceStatImpl;
51use common_telemetry::{info, warn};
52use snafu::{ResultExt, ensure};
53use store_api::storage::MAX_REGION_SEQ;
54
55use crate::bootstrap::build_default_meta_peer_client;
56use crate::cache_invalidator::MetasrvCacheInvalidator;
57use crate::cluster::MetaPeerClientRef;
58use crate::error::{self, BuildWalProviderSnafu, OtherSnafu, Result};
59use crate::event::EventHandlerImpl;
60use crate::gc::{DefaultGcSchedulerCtx, GcScheduler};
61use crate::greptimedb_telemetry::get_greptimedb_telemetry_task;
62use crate::handler::failure_handler::RegionFailureHandler;
63use crate::handler::flow_state_handler::FlowStateHandler;
64use crate::handler::persist_stats_handler::PersistStatsHandler;
65use crate::handler::region_lease_handler::{CustomizedRegionLeaseRenewerRef, RegionLeaseHandler};
66use crate::handler::{HeartbeatHandlerGroupBuilder, HeartbeatMailbox, Pushers};
67use crate::metasrv::{
68    ElectionRef, FLOW_ID_SEQ, METASRV_DATA_DIR, Metasrv, MetasrvInfo, MetasrvOptions,
69    RegionStatAwareSelectorRef, SelectTarget, SelectorContext, SelectorRef, TABLE_ID_SEQ,
70};
71use crate::peer::MetasrvPeerAllocator;
72use crate::procedure::region_migration::DefaultContextFactory;
73use crate::procedure::region_migration::manager::RegionMigrationManager;
74use crate::procedure::repartition::gc_requirement::RepartitionGcRequirementManager;
75use crate::procedure::repartition::{
76    DefaultRepartitionProcedureFactory, GcDisabledRepartitionProcedureFactory,
77};
78use crate::procedure::wal_prune::Context as WalPruneContext;
79use crate::procedure::wal_prune::manager::{WalPruneManager, WalPruneTicker};
80use crate::region::flush_trigger::RegionFlushTrigger;
81use crate::region::supervisor::{
82    DEFAULT_INITIALIZATION_RETRY_PERIOD, DEFAULT_TICK_INTERVAL, HeartbeatAcceptor,
83    RegionFailureDetectorControl, RegionSupervisor, RegionSupervisorSelector,
84    RegionSupervisorTicker,
85};
86use crate::selector::lease_based::LeaseBasedSelector;
87use crate::selector::round_robin::RoundRobinSelector;
88use crate::service::mailbox::MailboxRef;
89use crate::service::store::cached_kv::LeaderCachedKvBackend;
90use crate::state::State;
91use crate::utils::database::DatabaseOperator;
92use crate::utils::insert_forwarder::InsertForwarder;
93
94/// The time window for twcs compaction of the region stats table.
95const REGION_STATS_TABLE_TWCS_COMPACTION_TIME_WINDOW: Duration = Duration::from_secs(86400);
96
97// TODO(fys): try use derive_builder macro
98pub struct MetasrvBuilder {
99    options: Option<MetasrvOptions>,
100    kv_backend: Option<KvBackendRef>,
101    in_memory: Option<ResettableKvBackendRef>,
102    selector: Option<SelectorRef>,
103    handler_group_builder: Option<HeartbeatHandlerGroupBuilder>,
104    election: Option<ElectionRef>,
105    meta_peer_client: Option<MetaPeerClientRef>,
106    node_manager: Option<NodeManagerRef>,
107    plugins: Option<Plugins>,
108    table_metadata_allocator: Option<TableMetadataAllocatorRef>,
109}
110
111impl MetasrvBuilder {
112    pub fn new() -> Self {
113        Self {
114            kv_backend: None,
115            in_memory: None,
116            selector: None,
117            handler_group_builder: None,
118            meta_peer_client: None,
119            election: None,
120            options: None,
121            node_manager: None,
122            plugins: None,
123            table_metadata_allocator: None,
124        }
125    }
126
127    pub fn options(mut self, options: MetasrvOptions) -> Self {
128        self.options = Some(options);
129        self
130    }
131
132    pub fn kv_backend(mut self, kv_backend: KvBackendRef) -> Self {
133        self.kv_backend = Some(kv_backend);
134        self
135    }
136
137    pub fn in_memory(mut self, in_memory: ResettableKvBackendRef) -> Self {
138        self.in_memory = Some(in_memory);
139        self
140    }
141
142    pub fn selector(mut self, selector: SelectorRef) -> Self {
143        self.selector = Some(selector);
144        self
145    }
146
147    pub fn heartbeat_handler(
148        mut self,
149        handler_group_builder: HeartbeatHandlerGroupBuilder,
150    ) -> Self {
151        self.handler_group_builder = Some(handler_group_builder);
152        self
153    }
154
155    pub fn meta_peer_client(mut self, meta_peer_client: MetaPeerClientRef) -> Self {
156        self.meta_peer_client = Some(meta_peer_client);
157        self
158    }
159
160    pub fn election(mut self, election: Option<ElectionRef>) -> Self {
161        self.election = election;
162        self
163    }
164
165    pub fn node_manager(mut self, node_manager: NodeManagerRef) -> Self {
166        self.node_manager = Some(node_manager);
167        self
168    }
169
170    pub fn plugins(mut self, plugins: Plugins) -> Self {
171        self.plugins = Some(plugins);
172        self
173    }
174
175    pub fn table_metadata_allocator(
176        mut self,
177        table_metadata_allocator: TableMetadataAllocatorRef,
178    ) -> Self {
179        self.table_metadata_allocator = Some(table_metadata_allocator);
180        self
181    }
182
183    pub fn options_ref(&self) -> Option<&MetasrvOptions> {
184        self.options.as_ref()
185    }
186
187    pub fn kv_backend_ref(&self) -> Option<&KvBackendRef> {
188        self.kv_backend.as_ref()
189    }
190
191    pub fn in_memory_ref(&self) -> Option<&ResettableKvBackendRef> {
192        self.in_memory.as_ref()
193    }
194
195    pub fn election_ref(&self) -> Option<&ElectionRef> {
196        self.election.as_ref()
197    }
198
199    pub fn meta_peer_client_ref(&self) -> Option<&MetaPeerClientRef> {
200        self.meta_peer_client.as_ref()
201    }
202
203    pub fn node_manager_ref(&self) -> Option<&NodeManagerRef> {
204        self.node_manager.as_ref()
205    }
206
207    pub async fn build(self) -> Result<Metasrv> {
208        let MetasrvBuilder {
209            election,
210            meta_peer_client,
211            options,
212            kv_backend,
213            in_memory,
214            selector,
215            handler_group_builder,
216            node_manager,
217            plugins,
218            table_metadata_allocator,
219        } = self;
220
221        let options = options.unwrap_or_default();
222        options.gc.validate()?;
223
224        let kv_backend = kv_backend.unwrap_or_else(|| Arc::new(MemoryKvBackend::new()));
225        let in_memory = in_memory.unwrap_or_else(|| Arc::new(MemoryKvBackend::new()));
226
227        let state = Arc::new(RwLock::new(match election {
228            None => State::leader(options.grpc.server_addr.clone(), true),
229            Some(_) => State::follower(options.grpc.server_addr.clone()),
230        }));
231
232        let leader_cached_kv_backend = Arc::new(LeaderCachedKvBackend::new(
233            state.clone(),
234            kv_backend.clone(),
235        ));
236
237        let meta_peer_client = meta_peer_client
238            .unwrap_or_else(|| build_default_meta_peer_client(&election, &in_memory));
239        let database_operator = Arc::new(DatabaseOperator::new(meta_peer_client.clone()));
240
241        let event_inserter = Box::new(InsertForwarder::new(
242            database_operator.clone(),
243            Some(InsertOptions {
244                ttl: options.event_recorder.ttl,
245                append_mode: true,
246                twcs_compaction_time_window: Some(DEFAULT_COMPACTION_TIME_WINDOW),
247            }),
248        ));
249        // Builds the event recorder to record important events and persist them as the system table.
250        let event_recorder = Arc::new(EventRecorderImpl::with_event_type_filter(
251            Box::new(EventHandlerImpl::new(event_inserter)),
252            options.event_recorder.event_types.clone(),
253            options.event_recorder.flush_interval,
254        ));
255
256        let selector = selector.unwrap_or_else(|| Arc::new(LeaseBasedSelector));
257        let pushers = Pushers::default();
258        let mailbox = build_mailbox(&kv_backend, &pushers);
259        let runtime_switch_manager = Arc::new(RuntimeSwitchManager::new(kv_backend.clone()));
260        let repartition_gc_requirement_manager =
261            Arc::new(RepartitionGcRequirementManager::new(kv_backend.clone()));
262        let procedure_manager = build_procedure_manager(
263            &options,
264            &kv_backend,
265            &runtime_switch_manager,
266            event_recorder,
267        );
268
269        let table_metadata_manager = Arc::new(TableMetadataManager::new(
270            leader_cached_kv_backend.clone() as _,
271        ));
272        let flow_metadata_manager = Arc::new(FlowMetadataManager::new(
273            leader_cached_kv_backend.clone() as _,
274        ));
275
276        let selector_ctx = SelectorContext {
277            peer_discovery: meta_peer_client.clone(),
278        };
279
280        let wal_provider = build_wal_provider(&options.wal, kv_backend.clone())
281            .await
282            .context(BuildWalProviderSnafu)?;
283        let wal_provider = Arc::new(wal_provider);
284        let is_remote_wal = wal_provider.is_remote_wal();
285        let table_metadata_allocator = table_metadata_allocator.unwrap_or_else(|| {
286            let sequence = Arc::new(
287                SequenceBuilder::new(TABLE_ID_SEQ, kv_backend.clone())
288                    .initial(MIN_USER_TABLE_ID as u64)
289                    .step(10)
290                    .build(),
291            );
292            let peer_allocator = Arc::new(
293                MetasrvPeerAllocator::new(selector_ctx.clone(), selector.clone())
294                    .with_max_items(MAX_REGION_SEQ),
295            );
296            Arc::new(TableMetadataAllocator::with_peer_allocator(
297                sequence,
298                wal_provider.clone(),
299                peer_allocator,
300            ))
301        });
302        let table_id_allocator = table_metadata_allocator.table_id_allocator();
303
304        let flow_selector =
305            Arc::new(RoundRobinSelector::new(SelectTarget::Flownode)) as SelectorRef;
306
307        let flow_metadata_allocator = {
308            // for now flownode just use round-robin selector
309            let flow_selector_ctx = selector_ctx.clone();
310            let peer_allocator = Arc::new(MetasrvPeerAllocator::new(
311                flow_selector_ctx,
312                flow_selector.clone(),
313            ));
314            let seq = Arc::new(
315                SequenceBuilder::new(FLOW_ID_SEQ, kv_backend.clone())
316                    .initial(MIN_USER_FLOW_ID as u64)
317                    .step(10)
318                    .build(),
319            );
320
321            Arc::new(FlowMetadataAllocator::with_peer_allocator(
322                seq,
323                peer_allocator,
324            ))
325        };
326        let flow_state_handler =
327            FlowStateHandler::new(FlowStateManager::new(in_memory.clone().as_kv_backend_ref()));
328
329        let memory_region_keeper = Arc::new(MemoryRegionKeeper::default());
330        let node_manager = node_manager.unwrap_or_else(|| {
331            let datanode_client_channel_config = options.datanode.client.channel_config();
332            Arc::new(NodeClients::new(datanode_client_channel_config))
333        });
334        let cache_invalidator = Arc::new(MetasrvCacheInvalidator::new(
335            mailbox.clone(),
336            MetasrvInfo {
337                server_addr: options.grpc.server_addr.clone(),
338            },
339        ));
340
341        if !is_remote_wal && options.enable_region_failover {
342            ensure!(
343                options.allow_region_failover_on_local_wal,
344                error::UnexpectedSnafu {
345                    violated: "Region failover is not supported in the local WAL implementation!
346                    If you want to enable region failover for local WAL, please set `allow_region_failover_on_local_wal` to true.",
347                }
348            );
349            if options.allow_region_failover_on_local_wal {
350                warn!(
351                    "Region failover is force enabled in the local WAL implementation! This may lead to data loss during failover!"
352                );
353            }
354        }
355
356        let (tx, rx) = RegionSupervisor::channel();
357        let (region_failure_detector_controller, region_supervisor_ticker): (
358            RegionFailureDetectorControllerRef,
359            Option<std::sync::Arc<RegionSupervisorTicker>>,
360        ) = if options.enable_region_failover {
361            (
362                Arc::new(RegionFailureDetectorControl::new(tx.clone())) as _,
363                Some(Arc::new(RegionSupervisorTicker::new(
364                    DEFAULT_TICK_INTERVAL,
365                    options.region_failure_detector_initialization_delay,
366                    DEFAULT_INITIALIZATION_RETRY_PERIOD,
367                    tx.clone(),
368                ))),
369            )
370        } else {
371            (Arc::new(NoopRegionFailureDetectorControl) as _, None as _)
372        };
373
374        // region migration manager
375        let region_migration_manager = Arc::new(RegionMigrationManager::new(
376            procedure_manager.clone(),
377            DefaultContextFactory::new(
378                in_memory.clone(),
379                table_metadata_manager.clone(),
380                memory_region_keeper.clone(),
381                region_failure_detector_controller.clone(),
382                mailbox.clone(),
383                options.grpc.server_addr.clone(),
384                cache_invalidator.clone(),
385            ),
386        ));
387        region_migration_manager.try_start()?;
388        let region_supervisor_selector = plugins
389            .as_ref()
390            .and_then(|plugins| plugins.get::<RegionStatAwareSelectorRef>());
391
392        let supervisor_selector = match region_supervisor_selector {
393            Some(selector) => {
394                info!("Using region stat aware selector");
395                RegionSupervisorSelector::RegionStatAwareSelector(selector)
396            }
397            None => RegionSupervisorSelector::NaiveSelector(selector.clone()),
398        };
399
400        let region_failover_handler = if options.enable_region_failover {
401            let region_supervisor = RegionSupervisor::new(
402                rx,
403                options.failure_detector,
404                selector_ctx.clone(),
405                supervisor_selector,
406                region_migration_manager.clone(),
407                runtime_switch_manager.clone(),
408                meta_peer_client.clone(),
409                leader_cached_kv_backend.clone(),
410            )
411            .with_state(state.clone());
412
413            Some(RegionFailureHandler::new(
414                region_supervisor,
415                HeartbeatAcceptor::new(tx),
416            ))
417        } else {
418            None
419        };
420
421        let leader_region_registry = Arc::new(LeaderRegionRegistry::default());
422        let topic_stats_registry = Arc::new(TopicStatsRegistry::default());
423
424        let ddl_context = DdlContext {
425            node_manager: node_manager.clone(),
426            cache_invalidator: cache_invalidator.clone(),
427            memory_region_keeper: memory_region_keeper.clone(),
428            leader_region_registry: leader_region_registry.clone(),
429            table_metadata_manager: table_metadata_manager.clone(),
430            table_metadata_allocator: table_metadata_allocator.clone(),
431            flow_metadata_manager: flow_metadata_manager.clone(),
432            flow_metadata_allocator: flow_metadata_allocator.clone(),
433            region_failure_detector_controller,
434            soft_drop_enabled: ddl_soft_drop_enabled(&options),
435            soft_drop_retention: ddl_soft_drop_retention(&options),
436            create_database_metadata_committer: None,
437        };
438        let procedure_manager_c = procedure_manager.clone();
439        let repartition_procedure_factory: RepartitionProcedureFactoryRef = if options.gc.enable {
440            Arc::new(DefaultRepartitionProcedureFactory::new(
441                mailbox.clone(),
442                options.grpc.server_addr.clone(),
443                repartition_gc_requirement_manager.clone(),
444            ))
445        } else {
446            Arc::new(GcDisabledRepartitionProcedureFactory::new(
447                mailbox.clone(),
448                options.grpc.server_addr.clone(),
449                repartition_gc_requirement_manager.clone(),
450            ))
451        };
452        let ddl_manager = DdlManager::new(
453            ddl_context,
454            procedure_manager_c,
455            repartition_procedure_factory,
456        );
457
458        let ddl_manager = if let Some(configurator) = plugins
459            .as_ref()
460            .and_then(|p| p.get::<DdlManagerConfiguratorRef<DdlManagerConfigureContext>>())
461        {
462            let ctx = DdlManagerConfigureContext {
463                kv_backend: kv_backend.clone(),
464                meta_peer_client: meta_peer_client.clone(),
465            };
466            configurator
467                .configure(ddl_manager, ctx)
468                .await
469                .context(OtherSnafu)?
470        } else {
471            ddl_manager
472        };
473        ddl_manager
474            .register_loaders()
475            .context(error::InitDdlManagerSnafu)?;
476
477        let ddl_manager = Arc::new(ddl_manager);
478
479        let region_flush_ticker = if is_remote_wal {
480            let remote_wal_options = options.wal.remote_wal_options().unwrap();
481            let (region_flush_trigger, region_flush_ticker) = RegionFlushTrigger::new(
482                table_metadata_manager.clone(),
483                leader_region_registry.clone(),
484                topic_stats_registry.clone(),
485                mailbox.clone(),
486                options.grpc.server_addr.clone(),
487                remote_wal_options.flush_trigger_size,
488                remote_wal_options.checkpoint_trigger_size,
489                remote_wal_options.region_flush_trigger_interval,
490                remote_wal_options.periodic_checkpoint_persist_interval,
491            );
492            region_flush_trigger.try_start()?;
493
494            Some(Arc::new(region_flush_ticker))
495        } else {
496            None
497        };
498
499        // remote WAL prune ticker and manager
500        let wal_prune_ticker = if is_remote_wal && options.wal.enable_active_wal_pruning() {
501            let (tx, rx) = WalPruneManager::channel();
502            // Safety: Must be remote WAL.
503            let remote_wal_options = options.wal.remote_wal_options().unwrap();
504            let kafka_client = build_kafka_client(&remote_wal_options.connection)
505                .await
506                .context(error::BuildKafkaClientSnafu)?;
507            let wal_prune_context = WalPruneContext {
508                client: Arc::new(kafka_client),
509                table_metadata_manager: table_metadata_manager.clone(),
510                leader_region_registry: leader_region_registry.clone(),
511            };
512            let wal_prune_manager = WalPruneManager::new(
513                remote_wal_options.auto_prune_parallelism,
514                remote_wal_options.auto_prune_logical_delete,
515                rx,
516                procedure_manager.clone(),
517                wal_prune_context,
518            );
519            // Start manager in background. Ticker will be started in the main thread to send ticks.
520            wal_prune_manager.try_start().await?;
521            let wal_prune_ticker = Arc::new(WalPruneTicker::new(
522                remote_wal_options.auto_prune_interval,
523                tx.clone(),
524            ));
525            Some(wal_prune_ticker)
526        } else {
527            None
528        };
529
530        let gc_ticker = if options.gc.enable {
531            let gc_scheduler_ctx = DefaultGcSchedulerCtx::try_new(
532                table_metadata_manager.clone(),
533                procedure_manager.clone(),
534                runtime_switch_manager.clone(),
535                #[cfg(feature = "enterprise")]
536                ddl_manager.clone(),
537                meta_peer_client.clone(),
538                mailbox.clone(),
539                options.grpc.server_addr.clone(),
540            )?;
541            let (gc_scheduler, gc_ticker) = GcScheduler::new_with_config(
542                gc_scheduler_ctx,
543                runtime_switch_manager.clone(),
544                options.gc.clone(),
545            )?;
546            gc_scheduler.try_start()?;
547
548            Some(Arc::new(gc_ticker))
549        } else {
550            None
551        };
552
553        let customized_region_lease_renewer = plugins
554            .as_ref()
555            .and_then(|plugins| plugins.get::<CustomizedRegionLeaseRenewerRef>());
556
557        let persist_region_stats_handler = if !options.stats_persistence.ttl.is_zero() {
558            let inserter = Box::new(InsertForwarder::new(
559                database_operator.clone(),
560                Some(InsertOptions {
561                    ttl: options.stats_persistence.ttl,
562                    append_mode: true,
563                    twcs_compaction_time_window: Some(
564                        REGION_STATS_TABLE_TWCS_COMPACTION_TIME_WINDOW,
565                    ),
566                }),
567            ));
568
569            Some(PersistStatsHandler::new(
570                inserter,
571                options.stats_persistence.interval,
572            ))
573        } else {
574            None
575        };
576
577        let handler_group_builder = match handler_group_builder {
578            Some(handler_group_builder) => handler_group_builder,
579            None => {
580                let region_lease_handler = RegionLeaseHandler::new(
581                    default_distributed_time_constants().region_lease.as_secs(),
582                    table_metadata_manager.clone(),
583                    memory_region_keeper.clone(),
584                    customized_region_lease_renewer,
585                );
586
587                HeartbeatHandlerGroupBuilder::new(pushers)
588                    .with_plugins(plugins.clone())
589                    .with_region_failure_handler(region_failover_handler)
590                    .with_region_lease_handler(Some(region_lease_handler))
591                    .with_flush_stats_factor(Some(options.flush_stats_factor))
592                    .with_flow_state_handler(Some(flow_state_handler))
593                    .with_persist_stats_handler(persist_region_stats_handler)
594                    .add_default_handlers()
595            }
596        };
597
598        let enable_telemetry = options.enable_telemetry;
599        let metasrv_home = Path::new(&options.data_home)
600            .join(METASRV_DATA_DIR)
601            .to_string_lossy()
602            .to_string();
603
604        let reconciliation_manager = Arc::new(ReconciliationManager::new(
605            node_manager.clone(),
606            table_metadata_manager.clone(),
607            cache_invalidator.clone(),
608            procedure_manager.clone(),
609        ));
610        reconciliation_manager
611            .try_start()
612            .context(error::InitReconciliationManagerSnafu)?;
613
614        let mut resource_stat = ResourceStatImpl::default();
615        resource_stat.start_collect_cpu_usage();
616
617        Ok(Metasrv {
618            state,
619            started: Arc::new(AtomicBool::new(false)),
620            start_time_ms: common_time::util::current_time_millis() as u64,
621            options,
622            in_memory,
623            kv_backend,
624            leader_cached_kv_backend,
625            meta_peer_client: meta_peer_client.clone(),
626            selector,
627            selector_ctx,
628            // TODO(jeremy): We do not allow configuring the flow selector.
629            flow_selector,
630            handler_group: RwLock::new(None),
631            handler_group_builder: Mutex::new(Some(handler_group_builder)),
632            election,
633            procedure_manager,
634            mailbox,
635            ddl_manager,
636            wal_provider,
637            table_metadata_manager,
638            runtime_switch_manager,
639            repartition_gc_requirement_manager,
640            greptimedb_telemetry_task: get_greptimedb_telemetry_task(
641                Some(metasrv_home),
642                meta_peer_client,
643                enable_telemetry,
644            )
645            .await,
646            plugins: plugins.unwrap_or_else(Plugins::default),
647            memory_region_keeper,
648            region_migration_manager,
649            region_supervisor_ticker,
650            cache_invalidator,
651            leader_region_registry,
652            wal_prune_ticker,
653            region_flush_ticker,
654            table_id_allocator,
655            reconciliation_manager,
656            topic_stats_registry,
657            resource_stat: Arc::new(resource_stat),
658            gc_ticker,
659            database_operator,
660        })
661    }
662}
663
664fn build_mailbox(kv_backend: &KvBackendRef, pushers: &Pushers) -> MailboxRef {
665    let mailbox_sequence = SequenceBuilder::new("heartbeat_mailbox", kv_backend.clone())
666        .initial(1)
667        .step(100)
668        .build();
669
670    HeartbeatMailbox::create(pushers.clone(), mailbox_sequence)
671}
672
673fn build_procedure_manager(
674    options: &MetasrvOptions,
675    kv_backend: &KvBackendRef,
676    runtime_switch_manager: &RuntimeSwitchManagerRef,
677    event_recorder: EventRecorderRef,
678) -> ProcedureManagerRef {
679    let manager_config = ManagerConfig {
680        max_retry_times: options.procedure.max_retry_times,
681        retry_delay: options.procedure.retry_delay,
682        max_running_procedures: options.procedure.max_running_procedures,
683        ..Default::default()
684    };
685    let kv_state_store = Arc::new(
686        KvStateStore::new(kv_backend.clone()).with_max_value_size(
687            options
688                .procedure
689                .max_metadata_value_size
690                .map(|v| v.as_bytes() as usize),
691        ),
692    );
693
694    Arc::new(LocalManager::new(
695        manager_config,
696        kv_state_store.clone(),
697        kv_state_store,
698        Some(runtime_switch_manager.clone()),
699        Some(event_recorder),
700    ))
701}
702
703/// Resolves if soft-drop is enabled from metasrv options.
704///
705/// Soft drop is an enterprise-only feature; it is always disabled in
706/// non-enterprise builds regardless of the configuration.
707fn ddl_soft_drop_enabled(options: &MetasrvOptions) -> bool {
708    cfg!(feature = "enterprise") && options.gc.experimental_soft_drop.enable
709}
710
711/// Returns soft-drop retention for recovering persisted procedures.
712fn ddl_soft_drop_retention(options: &MetasrvOptions) -> Option<Duration> {
713    Some(options.gc.experimental_soft_drop.retention)
714}
715
716impl Default for MetasrvBuilder {
717    fn default() -> Self {
718        Self::new()
719    }
720}
721
722/// The context for [`DdlManagerConfiguratorRef`].
723pub struct DdlManagerConfigureContext {
724    pub kv_backend: KvBackendRef,
725    pub meta_peer_client: MetaPeerClientRef,
726}
727
728#[cfg(test)]
729mod tests {
730    use super::*;
731
732    #[cfg(feature = "enterprise")]
733    #[test]
734    fn test_ddl_soft_drop_gate_preserves_retention_for_recovery() {
735        let mut options = MetasrvOptions::default();
736        options.gc.enable = true;
737        options.gc.experimental_soft_drop.enable = true;
738        options.gc.experimental_soft_drop.retention = Duration::from_secs(123);
739
740        assert!(ddl_soft_drop_enabled(&options));
741        assert_eq!(
742            Some(Duration::from_secs(123)),
743            ddl_soft_drop_retention(&options)
744        );
745    }
746
747    #[cfg(not(feature = "enterprise"))]
748    #[test]
749    fn test_ddl_soft_drop_is_always_disabled_in_non_enterprise_build() {
750        let mut options = MetasrvOptions::default();
751        options.gc.enable = true;
752        options.gc.experimental_soft_drop.enable = true;
753        options.gc.experimental_soft_drop.retention = Duration::from_secs(123);
754
755        assert!(!ddl_soft_drop_enabled(&options));
756        // Retention is preserved for recovering persisted procedures.
757        assert_eq!(
758            Some(Duration::from_secs(123)),
759            ddl_soft_drop_retention(&options)
760        );
761    }
762
763    #[test]
764    fn test_ddl_soft_drop_is_disabled_by_default() {
765        assert!(!ddl_soft_drop_enabled(&MetasrvOptions::default()));
766    }
767
768    #[test]
769    fn test_soft_drop_options_are_validated() {
770        let mut options = MetasrvOptions::default();
771        options.gc.enable = true;
772        options.gc.experimental_soft_drop.enable = true;
773        options.gc.experimental_soft_drop.retention = Duration::ZERO;
774
775        assert!(options.gc.validate().is_err());
776    }
777
778    #[tokio::test]
779    async fn test_builder_rejects_soft_drop_when_gc_is_disabled() {
780        let mut options = MetasrvOptions::default();
781        options.gc.experimental_soft_drop.enable = true;
782
783        assert!(
784            MetasrvBuilder::new()
785                .options(options)
786                .build()
787                .await
788                .is_err()
789        );
790    }
791
792    #[tokio::test]
793    async fn test_builder_skips_gc_validation_when_gc_is_disabled() {
794        let mut options = MetasrvOptions::default();
795        options.gc.enable = false;
796        options.gc.max_concurrent_tables = 0;
797
798        MetasrvBuilder::new()
799            .options(options)
800            .build()
801            .await
802            .unwrap();
803    }
804}