mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-07 13:53:12 +00:00
31c6859965
Second TODO-convergence round over the current tree (backlog#646). All line numbers in the old inventory had gone stale after the set_disk / diagnostics / cluster refactors, so this re-scans and reduces the marker count from 144 to 99. STALE removals (comment describes already-implemented behavior, or dead commented-out blocks) across ecstore (set_disk ops/core, store, cluster/rpc, bucket/metadata_sys, services), iam, filemeta, s3select and rustfs auth/object_usecase. No behavior change. Safe fills, each verified: - filemeta: replication_info_equals now also compares replication_state_internal (function currently has no callers; adds a regression test). - bitrot: drop the confirmed-unused `_want` parameter from bitrot_verify and the now-unused `sum` on LocalDisk::bitrot_verify, removing a Bytes::copy_from_slice allocation. Streaming verify uses the file's embedded per-shard hash, never the passed sum. - signer: rename v4_ignored_headers -> V4_IGNORED_HEADERS and drop the non_upper_case_globals allow. - admin/heal: test_decode was #[ignore]d and used serde_urlencoded on a JSON body (would panic); rewire to serde_json::from_slice to match the production decode path, add assertions, un-ignore. Verified: cargo fmt; cargo check on touched crates; tests pass (filemeta, signer, bitrot, heal::test_decode); arch guardrail scripts pass.
145 lines
4.5 KiB
Rust
145 lines
4.5 KiB
Rust
// Copyright 2024 RustFS Team
|
|
//
|
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
// you may not use this file except in compliance with the License.
|
|
// You may obtain a copy of the License at
|
|
//
|
|
// http://www.apache.org/licenses/LICENSE-2.0
|
|
//
|
|
// Unless required by applicable law or agreed to in writing, software
|
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
// See the License for the specific language governing permissions and
|
|
// limitations under the License.
|
|
|
|
use std::sync::Arc;
|
|
|
|
use async_trait::async_trait;
|
|
use datafusion::arrow::datatypes::DataType;
|
|
use datafusion::common::{Result as DFResult, TableReference};
|
|
use datafusion::datasource::TableProvider;
|
|
use datafusion::logical_expr::var_provider::is_system_variables;
|
|
use datafusion::logical_expr::{AggregateUDF, HigherOrderUDF, ScalarUDF, TableSource, WindowUDF};
|
|
use datafusion::variable::VarType;
|
|
use datafusion::{config::ConfigOptions, sql::planner::ContextProvider};
|
|
use rustfs_s3select_api::query::{function::FuncMetaManagerRef, session::SessionCtx};
|
|
|
|
use crate::data_source::table_source::{TableHandle, TableSourceAdapter};
|
|
|
|
pub mod base_table;
|
|
|
|
#[async_trait]
|
|
pub trait ContextProviderExtension: ContextProvider {
|
|
fn get_table_source_(&self, name: TableReference) -> datafusion::common::Result<Arc<TableSourceAdapter>>;
|
|
}
|
|
|
|
pub type TableHandleProviderRef = Arc<dyn TableHandleProvider + Send + Sync>;
|
|
|
|
pub trait TableHandleProvider {
|
|
fn build_table_handle(&self, provider: Arc<dyn TableProvider>) -> DFResult<TableHandle>;
|
|
}
|
|
|
|
pub struct MetadataProvider {
|
|
provider: Arc<dyn TableProvider>,
|
|
session: SessionCtx,
|
|
config_options: ConfigOptions,
|
|
func_manager: FuncMetaManagerRef,
|
|
current_session_table_provider: TableHandleProviderRef,
|
|
}
|
|
|
|
impl MetadataProvider {
|
|
#[allow(clippy::too_many_arguments)]
|
|
pub fn new(
|
|
provider: Arc<dyn TableProvider>,
|
|
current_session_table_provider: TableHandleProviderRef,
|
|
func_manager: FuncMetaManagerRef,
|
|
session: SessionCtx,
|
|
) -> Self {
|
|
Self {
|
|
provider,
|
|
current_session_table_provider,
|
|
config_options: session.inner().config_options().as_ref().clone(),
|
|
session,
|
|
func_manager,
|
|
}
|
|
}
|
|
|
|
fn build_table_handle(&self) -> datafusion::common::Result<TableHandle> {
|
|
self.current_session_table_provider.build_table_handle(self.provider.clone())
|
|
}
|
|
}
|
|
|
|
impl ContextProviderExtension for MetadataProvider {
|
|
fn get_table_source_(&self, table_ref: TableReference) -> datafusion::common::Result<Arc<TableSourceAdapter>> {
|
|
let name = table_ref.clone().resolve("", "");
|
|
let table_name = &*name.table;
|
|
|
|
let table_handle = self.build_table_handle()?;
|
|
|
|
Ok(Arc::new(TableSourceAdapter::try_new(table_ref.clone(), table_name, table_handle)?))
|
|
}
|
|
}
|
|
|
|
impl ContextProvider for MetadataProvider {
|
|
fn get_function_meta(&self, name: &str) -> Option<Arc<ScalarUDF>> {
|
|
self.func_manager
|
|
.udf(name)
|
|
.ok()
|
|
.or_else(|| self.session.inner().scalar_functions().get(name).cloned())
|
|
}
|
|
|
|
fn get_aggregate_meta(&self, name: &str) -> Option<Arc<AggregateUDF>> {
|
|
self.func_manager.udaf(name).ok()
|
|
}
|
|
|
|
fn get_higher_order_meta(&self, _name: &str) -> Option<Arc<HigherOrderUDF>> {
|
|
None
|
|
}
|
|
|
|
fn get_variable_type(&self, variable_names: &[String]) -> Option<DataType> {
|
|
if variable_names.is_empty() {
|
|
return None;
|
|
}
|
|
|
|
let var_type = if is_system_variables(variable_names) {
|
|
VarType::System
|
|
} else {
|
|
VarType::UserDefined
|
|
};
|
|
|
|
self.session
|
|
.inner()
|
|
.execution_props()
|
|
.get_var_provider(var_type)
|
|
.and_then(|p| p.get_type(variable_names))
|
|
}
|
|
|
|
fn options(&self) -> &ConfigOptions {
|
|
&self.config_options
|
|
}
|
|
|
|
fn get_window_meta(&self, name: &str) -> Option<Arc<WindowUDF>> {
|
|
self.func_manager.udwf(name).ok()
|
|
}
|
|
|
|
fn get_table_source(&self, name: TableReference) -> DFResult<Arc<dyn TableSource>> {
|
|
Ok(self.get_table_source_(name)?)
|
|
}
|
|
|
|
fn udf_names(&self) -> Vec<String> {
|
|
self.func_manager.udfs()
|
|
}
|
|
|
|
fn higher_order_function_names(&self) -> Vec<String> {
|
|
Vec::new()
|
|
}
|
|
|
|
fn udaf_names(&self) -> Vec<String> {
|
|
self.func_manager.udafs()
|
|
}
|
|
|
|
fn udwf_names(&self) -> Vec<String> {
|
|
self.func_manager.udwfs()
|
|
}
|
|
}
|