Skip to main content

frontend/instance/
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::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
62/// The frontend [`Instance`] builder.
63pub 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        // Builds cache registry
113        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        // Inject the entity-graph provider now that the query engine exists, so the
276        // computed `greptime_private.semantic_entities` / `semantic_relationships`
277        // tables can derive rows. Late binding here breaks the `catalog -> query`
278        // dependency cycle.
279        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}