// 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 futures::stream::AbortHandle; use parking_lot::Mutex; use rustfs_s3select_api::query::execution::{Output, QueryExecution, QueryStateMachineRef}; use rustfs_s3select_api::query::logical_planner::QueryPlan; use rustfs_s3select_api::query::optimizer::Optimizer; use rustfs_s3select_api::query::scheduler::SchedulerRef; use rustfs_s3select_api::{QueryError, QueryResult}; use tracing::debug; pub struct SqlQueryExecution { query_state_machine: QueryStateMachineRef, plan: QueryPlan, optimizer: Arc, scheduler: SchedulerRef, abort_handle: Mutex>, } impl SqlQueryExecution { pub fn new( query_state_machine: QueryStateMachineRef, plan: QueryPlan, optimizer: Arc, scheduler: SchedulerRef, ) -> Self { Self { query_state_machine, plan, optimizer, scheduler, abort_handle: Mutex::new(None), } } async fn start(&self) -> QueryResult { let physical_plan = { // Time optimize phase - dropped at end of this block let _optimize_timer = self.query_state_machine.time_phase("optimize"); self.query_state_machine.begin_optimize(); let plan = self.optimizer.optimize(&self.plan, &self.query_state_machine.session).await?; self.query_state_machine.end_optimize(); plan }; let stream = { // Time schedule phase - dropped at end of this block let _schedule_timer = self.query_state_machine.time_phase("schedule"); self.query_state_machine.begin_schedule(); let stream = self .scheduler .schedule(physical_plan.clone(), self.query_state_machine.session.inner().task_ctx()) .await? .stream(); self.query_state_machine.end_schedule(); stream }; Ok(Output::StreamData(stream)) } } #[async_trait] impl QueryExecution for SqlQueryExecution { async fn start(&self) -> QueryResult { let (task, abort_handle) = futures::future::abortable(self.start()); { *self.abort_handle.lock() = Some(abort_handle); } task.await.map_err(|_| QueryError::Cancel)? } fn cancel(&self) -> QueryResult<()> { debug!( "cancel sql query execution: sql: {}, state: {:?}", self.query_state_machine.query.content(), self.query_state_machine.state() ); // change state self.query_state_machine.cancel(); // stop future task if let Some(e) = self.abort_handle.lock().as_ref() { e.abort() }; debug!( "canceled sql query execution: sql: {}, state: {:?}", self.query_state_machine.query.content(), self.query_state_machine.state() ); Ok(()) } }