1use 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
94const REGION_STATS_TABLE_TWCS_COMPACTION_TIME_WINDOW: Duration = Duration::from_secs(86400);
96
97pub 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 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 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 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 let wal_prune_ticker = if is_remote_wal && options.wal.enable_active_wal_pruning() {
501 let (tx, rx) = WalPruneManager::channel();
502 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 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 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
703fn ddl_soft_drop_enabled(options: &MetasrvOptions) -> bool {
708 cfg!(feature = "enterprise") && options.gc.experimental_soft_drop.enable
709}
710
711fn 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
722pub 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 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}