Plan SQL via DataFusion and dispatch tasks to the Data Plane.
(
ctx: &DispatchCtx<'_>,
seq: u64,
sql: &str,
database_id: crate::types::DatabaseId,
)
| 140 | |
| 141 | /// Plan SQL via DataFusion and dispatch tasks to the Data Plane. |
| 142 | async fn execute_planned( |
| 143 | ctx: &DispatchCtx<'_>, |
| 144 | seq: u64, |
| 145 | sql: &str, |
| 146 | database_id: crate::types::DatabaseId, |
| 147 | ) -> NativeResponse { |
| 148 | // Extract per-query ON DENY override (e.g., SELECT ... ON DENY ERROR 'CODE' MESSAGE '...'). |
| 149 | let mut auth_ctx = ctx.auth_context.clone(); |
| 150 | let clean_sql = |
| 151 | crate::control::server::session_auth::extract_and_apply_on_deny(sql, &mut auth_ctx); |
| 152 | |
| 153 | let perm_cache = ctx.state.permission_cache.read().await; |
| 154 | let sec = crate::control::planner::context::PlanSecurityContext { |
| 155 | identity: ctx.identity, |
| 156 | auth: &auth_ctx, |
| 157 | rls_store: &ctx.state.rls, |
| 158 | permissions: &ctx.state.permissions, |
| 159 | roles: &ctx.state.roles, |
| 160 | permission_cache: Some(&*perm_cache), |
| 161 | }; |
| 162 | let tasks = match ctx |
| 163 | .query_ctx |
| 164 | .plan_sql_with_rls(&clean_sql, ctx.tenant_id(), database_id, &sec) |
| 165 | .await |
| 166 | { |
| 167 | Ok(t) => t, |
| 168 | Err(e) => return error_to_native(seq, &e), |
| 169 | }; |
| 170 | |
| 171 | if tasks.is_empty() { |
| 172 | return NativeResponse::status_row(seq, "OK"); |
| 173 | } |
| 174 | |
| 175 | let mut all_columns: Option<Vec<String>> = None; |
| 176 | let mut all_rows: Vec<Vec<Value>> = Vec::new(); |
| 177 | let mut last_lsn = 0u64; |
| 178 | let mut total_affected = 0u64; |
| 179 | |
| 180 | for task in tasks { |
| 181 | if task.tenant_id != ctx.tenant_id() { |
| 182 | return NativeResponse::error(seq, "42501", "tenant isolation violation"); |
| 183 | } |
| 184 | |
| 185 | // In transaction: buffer writes. |
| 186 | if ctx.sessions.transaction_state(ctx.peer_addr) == TransactionState::InBlock { |
| 187 | let is_write = crate::control::wal_replication::to_replicated_entry( |
| 188 | task.tenant_id, |
| 189 | task.vshard_id, |
| 190 | &task.plan, |
| 191 | ) |
| 192 | .is_some(); |
| 193 | if is_write { |
| 194 | ctx.sessions.buffer_write(ctx.peer_addr, task); |
| 195 | total_affected += 1; |
| 196 | continue; |
| 197 | } |
| 198 | } |
| 199 |
no test coverage detected