1use std::sync::Arc;
16use std::sync::atomic::AtomicBool;
17
18use auth::PermissionCheckerRef;
19use cache::{PARTITION_INFO_CACHE_NAME, TABLE_FLOWNODE_SET_CACHE_NAME, TABLE_ROUTE_CACHE_NAME};
20use catalog::CatalogManagerRef;
21use catalog::kvbackend::KvBackendCatalogManager;
22use catalog::process_manager::ProcessManagerRef;
23use catalog::system_schema::semantic_graph::EntityGraphProviderRef;
24use common_base::Plugins;
25use common_datasource::object_store::LocalFileAccess;
26use common_event_recorder::{EventRecorderImpl, EventRecorderRef};
27use common_meta::cache::{LayeredCacheRegistryRef, TableRouteCacheRef};
28use common_meta::cache_invalidator::{CacheInvalidatorRef, DummyCacheInvalidator};
29use common_meta::key::TableMetadataManager;
30use common_meta::key::flow::FlowMetadataManager;
31use common_meta::kv_backend::KvBackendRef;
32use common_meta::node_manager::NodeManagerRef;
33use common_meta::procedure_executor::ProcedureExecutorRef;
34use dashmap::DashMap;
35use operator::delete::Deleter;
36use operator::flow::FlowServiceOperator;
37use operator::insert::Inserter;
38use operator::procedure::ProcedureServiceOperator;
39use operator::request::Requester;
40#[cfg(feature = "enterprise")]
41use operator::statement::CreateDatabaseHandlerRef;
42use operator::statement::{
43 AdminEventRecorderHandle, AdminFunctionRecordingLayer, ExecutorConfigureContext,
44 StatementExecutor, StatementExecutorConfiguratorRef, StatementExecutorRef,
45};
46use operator::table::TableMutationOperator;
47use partition::cache::PartitionInfoCacheRef;
48use partition::manager::PartitionRuleManager;
49use pipeline::pipeline_operator::PipelineOperator;
50use query::QueryEngineFactory;
51use query::region_query::RegionQueryHandlerFactoryRef;
52use snafu::{OptionExt, ResultExt};
53
54use crate::error::{self, DataFusionSnafu, ExternalSnafu, Result};
55use crate::events::EventHandlerImpl;
56use crate::frontend::FrontendOptions;
57use crate::heartbeat::frontend_peer_addr;
58use crate::instance::Instance;
59use crate::instance::entity_graph::EntityGraphProviderImpl;
60use crate::instance::region_query::FrontendRegionQueryHandler;
61
62pub struct FrontendBuilder {
64 options: FrontendOptions,
65 kv_backend: KvBackendRef,
66 layered_cache_registry: LayeredCacheRegistryRef,
67 local_cache_invalidator: Option<CacheInvalidatorRef>,
68 catalog_manager: CatalogManagerRef,
69 node_manager: NodeManagerRef,
70 plugins: Option<Plugins>,
71 procedure_executor: ProcedureExecutorRef,
72 process_manager: ProcessManagerRef,
73 local_file_access: LocalFileAccess,
74}
75
76impl FrontendBuilder {
77 #[allow(clippy::too_many_arguments)]
78 pub fn new(
79 options: FrontendOptions,
80 kv_backend: KvBackendRef,
81 layered_cache_registry: LayeredCacheRegistryRef,
82 catalog_manager: CatalogManagerRef,
83 node_manager: NodeManagerRef,
84 procedure_executor: ProcedureExecutorRef,
85 process_manager: ProcessManagerRef,
86 ) -> Self {
87 Self {
88 options,
89 kv_backend,
90 layered_cache_registry,
91 local_cache_invalidator: None,
92 catalog_manager,
93 node_manager,
94 plugins: None,
95 procedure_executor,
96 process_manager,
97 local_file_access: LocalFileAccess::Disabled,
98 }
99 }
100
101 #[cfg(test)]
102 pub(crate) fn new_test(
103 options: &FrontendOptions,
104 meta_client: meta_client::MetaClientRef,
105 ) -> Self {
106 use cache::{build_fundamental_cache_registry, with_default_composite_cache_registry};
107 use common_meta::cache::LayeredCacheRegistryBuilder;
108 use common_meta::kv_backend::memory::MemoryKvBackend;
109
110 let kv_backend = Arc::new(MemoryKvBackend::new());
111
112 let layered_cache_builder = LayeredCacheRegistryBuilder::default();
114 let fundamental_cache_registry = build_fundamental_cache_registry(kv_backend.clone());
115 let layered_cache_registry = Arc::new(
116 with_default_composite_cache_registry(
117 layered_cache_builder.add_cache_registry(fundamental_cache_registry),
118 )
119 .unwrap()
120 .build(),
121 );
122
123 Self::new(
124 options.clone(),
125 kv_backend,
126 layered_cache_registry,
127 catalog::memory::MemoryCatalogManager::with_default_setup(),
128 Arc::new(client::client_manager::NodeClients::default()),
129 meta_client,
130 Arc::new(catalog::process_manager::ProcessManager::new(
131 "".to_string(),
132 None,
133 )),
134 )
135 }
136
137 pub fn with_local_cache_invalidator(self, cache_invalidator: CacheInvalidatorRef) -> Self {
138 Self {
139 local_cache_invalidator: Some(cache_invalidator),
140 ..self
141 }
142 }
143
144 pub fn with_local_file_access(mut self, local_file_access: LocalFileAccess) -> Self {
145 self.local_file_access = local_file_access;
146 self
147 }
148
149 pub fn options(&self) -> &FrontendOptions {
150 &self.options
151 }
152
153 pub fn kv_backend(&self) -> &KvBackendRef {
154 &self.kv_backend
155 }
156
157 pub fn layered_cache_registry(&self) -> &LayeredCacheRegistryRef {
158 &self.layered_cache_registry
159 }
160
161 pub fn catalog_manager(&self) -> &CatalogManagerRef {
162 &self.catalog_manager
163 }
164
165 pub fn node_manager(&self) -> &NodeManagerRef {
166 &self.node_manager
167 }
168
169 pub fn procedure_executor(&self) -> &ProcedureExecutorRef {
170 &self.procedure_executor
171 }
172
173 pub fn process_manager(&self) -> &ProcessManagerRef {
174 &self.process_manager
175 }
176
177 pub fn with_plugin(self, plugins: Plugins) -> Self {
178 Self {
179 plugins: Some(plugins),
180 ..self
181 }
182 }
183
184 pub async fn try_build(self) -> Result<Instance> {
185 let kv_backend = self.kv_backend;
186 let node_manager = self.node_manager;
187 let plugins = self.plugins.unwrap_or_default();
188 let process_manager = self.process_manager;
189 let table_route_cache: TableRouteCacheRef =
190 self.layered_cache_registry
191 .get()
192 .context(error::CacheRequiredSnafu {
193 name: TABLE_ROUTE_CACHE_NAME,
194 })?;
195 let partition_info_cache: PartitionInfoCacheRef = self
196 .layered_cache_registry
197 .get()
198 .context(error::CacheRequiredSnafu {
199 name: PARTITION_INFO_CACHE_NAME,
200 })?;
201 let partition_manager = Arc::new(PartitionRuleManager::new(
202 kv_backend.clone(),
203 table_route_cache.clone(),
204 partition_info_cache.clone(),
205 ));
206
207 let local_cache_invalidator = self
208 .local_cache_invalidator
209 .unwrap_or_else(|| Arc::new(DummyCacheInvalidator));
210
211 let region_query_handler =
212 if let Some(factory) = plugins.get::<RegionQueryHandlerFactoryRef>() {
213 factory.build(partition_manager.clone(), node_manager.clone())
214 } else {
215 FrontendRegionQueryHandler::arc(partition_manager.clone(), node_manager.clone())
216 };
217
218 let table_flownode_cache =
219 self.layered_cache_registry
220 .get()
221 .context(error::CacheRequiredSnafu {
222 name: TABLE_FLOWNODE_SET_CACHE_NAME,
223 })?;
224
225 let inserter = Arc::new(Inserter::new(
226 self.catalog_manager.clone(),
227 partition_manager.clone(),
228 node_manager.clone(),
229 table_flownode_cache,
230 self.options.auto_create_table,
231 ));
232 let deleter = Arc::new(Deleter::new(
233 self.catalog_manager.clone(),
234 partition_manager.clone(),
235 node_manager.clone(),
236 ));
237 let requester = Arc::new(Requester::new(
238 self.catalog_manager.clone(),
239 partition_manager.clone(),
240 node_manager.clone(),
241 ));
242 let table_mutation_handler = Arc::new(TableMutationOperator::new(
243 inserter.clone(),
244 deleter.clone(),
245 requester,
246 ));
247
248 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
249 let procedure_service_handler = Arc::new(ProcedureServiceOperator::new(
250 self.procedure_executor.clone(),
251 self.catalog_manager.clone(),
252 table_metadata_manager.clone(),
253 ));
254
255 let flow_metadata_manager: Arc<FlowMetadataManager> =
256 Arc::new(FlowMetadataManager::new(kv_backend.clone()));
257 let flow_service = FlowServiceOperator::new(flow_metadata_manager, node_manager.clone());
258
259 let mut query_options = self.options.query.clone();
260 query_options.enable_per_region_metrics = self.options.logging.enable_per_region_metrics;
261 let query_engine = QueryEngineFactory::try_new_with_plugins(
262 self.catalog_manager.clone(),
263 Some(partition_manager.clone()),
264 Some(region_query_handler.clone()),
265 Some(table_mutation_handler),
266 Some(procedure_service_handler),
267 Some(Arc::new(flow_service)),
268 true,
269 plugins.clone(),
270 query_options,
271 )
272 .context(DataFusionSnafu)?
273 .query_engine();
274
275 if let Some(kv_catalog) = self
280 .catalog_manager
281 .as_any()
282 .downcast_ref::<KvBackendCatalogManager>()
283 {
284 let provider: EntityGraphProviderRef = Arc::new(EntityGraphProviderImpl::new(
285 query_engine.clone(),
286 Arc::downgrade(&self.catalog_manager),
287 plugins.get::<PermissionCheckerRef>(),
288 ));
289 kv_catalog.set_entity_graph_provider(provider);
290 }
291
292 let frontend_peer_addr = frontend_peer_addr(&self.options);
293 let statement_executor = StatementExecutor::new(
294 self.catalog_manager.clone(),
295 query_engine.clone(),
296 self.procedure_executor,
297 kv_backend.clone(),
298 local_cache_invalidator,
299 inserter.clone(),
300 partition_manager,
301 Some(process_manager.clone()),
302 frontend_peer_addr.clone(),
303 self.local_file_access,
304 );
305
306 let statement_executor =
307 if let Some(configurator) = plugins.get::<StatementExecutorConfiguratorRef>() {
308 let ctx = ExecutorConfigureContext {
309 kv_backend: kv_backend.clone(),
310 };
311 configurator
312 .configure(statement_executor, ctx)
313 .await
314 .context(ExternalSnafu)?
315 } else {
316 statement_executor
317 };
318
319 #[cfg(feature = "enterprise")]
320 let statement_executor = if let Some(handler) = plugins.get::<CreateDatabaseHandlerRef>() {
321 statement_executor.with_create_database_handler(handler)
322 } else {
323 statement_executor
324 };
325
326 let admin_event_recorder = AdminEventRecorderHandle::default();
327 let statement_executor = statement_executor.with_admin_function_layer(Arc::new(
328 AdminFunctionRecordingLayer::new(admin_event_recorder.clone()),
329 ));
330 let statement_executor = Arc::new(statement_executor);
331
332 let pipeline_operator = Arc::new(PipelineOperator::new(
333 inserter.clone(),
334 statement_executor.clone(),
335 self.catalog_manager.clone(),
336 query_engine.clone(),
337 &self.options.pipeline,
338 ));
339
340 plugins.insert::<StatementExecutorRef>(statement_executor.clone());
341
342 let slow_query_recorder = Arc::new(EventRecorderImpl::new(Box::new(
343 EventHandlerImpl::new(statement_executor.clone(), self.options.slow_query.ttl),
344 )));
345 let event_recorder: EventRecorderRef = Arc::new(EventRecorderImpl::with_event_type_filter(
346 Box::new(EventHandlerImpl::new(
347 statement_executor.clone(),
348 self.options.event_recorder.ttl,
349 )),
350 self.options.event_recorder.event_types.clone(),
351 ));
352 admin_event_recorder.install(&event_recorder);
353
354 Ok(Instance {
355 frontend_peer_addr,
356 catalog_manager: self.catalog_manager,
357 pipeline_operator,
358 statement_executor,
359 query_engine,
360 plugins,
361 inserter,
362 deleter,
363 table_metadata_manager,
364 event_recorder,
365 slow_query_recorder,
366 process_manager,
367 otlp_metrics_table_legacy_cache: DashMap::new(),
368 slow_query_options: self.options.slow_query.clone(),
369 influxdb_default_merge_mode: self.options.influxdb.default_merge_mode,
370 trace_ingest_chunk_size: self.options.otlp.trace_ingest_chunk_size,
371 otlp_resource_info: self.options.otlp.experimental_enable_resource_info,
372 suspend: Arc::new(AtomicBool::new(false)),
373 })
374 }
375}