1use std::collections::BTreeMap;
11
12use maplit::btreemap;
13use mz_catalog::memory::error::ErrorKind;
14use mz_catalog::memory::objects::{CatalogItem, Index};
15use mz_ore::instrument;
16use mz_repr::explain::{ExprHumanizerExt, TransientItem};
17use mz_repr::optimize::{OptimizerFeatures, OverrideFrom};
18use mz_repr::{Datum, Row};
19use mz_sql::ast::ExplainStage;
20use mz_sql::catalog::CatalogError;
21use mz_sql::names::ResolvedIds;
22use mz_sql::plan;
23use mz_sql::session::metadata::SessionMetadata;
24use tracing::Span;
25
26use crate::command::ExecuteResponse;
27use crate::coord::sequencer::inner::return_if_err;
28use crate::coord::{
29 Coordinator, CreateIndexExplain, CreateIndexFinish, CreateIndexOptimize, CreateIndexStage,
30 ExplainContext, ExplainPlanContext, Message, PlanValidity, StageResult, Staged,
31};
32use crate::error::AdapterError;
33use crate::explain::explain_dataflow;
34use crate::explain::optimizer_trace::OptimizerTrace;
35use crate::optimize::dataflows::dataflow_import_id_bundle;
36use crate::optimize::{self, Optimize};
37use crate::session::Session;
38use crate::{AdapterNotice, ExecuteContext, catalog};
39
40impl Staged for CreateIndexStage {
41 type Ctx = ExecuteContext;
42
43 fn validity(&mut self) -> &mut PlanValidity {
44 match self {
45 Self::Optimize(stage) => &mut stage.validity,
46 Self::Finish(stage) => &mut stage.validity,
47 Self::Explain(stage) => &mut stage.validity,
48 }
49 }
50
51 async fn stage(
52 self,
53 coord: &mut Coordinator,
54 ctx: &mut ExecuteContext,
55 ) -> Result<StageResult<Box<Self>>, AdapterError> {
56 match self {
57 CreateIndexStage::Optimize(stage) => coord.create_index_optimize(stage).await,
58 CreateIndexStage::Finish(stage) => coord.create_index_finish(ctx, stage).await,
59 CreateIndexStage::Explain(stage) => {
60 coord.create_index_explain(ctx.session(), stage).await
61 }
62 }
63 }
64
65 fn message(self, ctx: ExecuteContext, span: Span) -> Message {
66 Message::CreateIndexStageReady {
67 ctx,
68 span,
69 stage: self,
70 }
71 }
72
73 fn cancel_enabled(&self) -> bool {
74 true
75 }
76}
77
78impl Coordinator {
79 #[instrument]
80 pub(crate) async fn sequence_create_index(
81 &mut self,
82 ctx: ExecuteContext,
83 plan: plan::CreateIndexPlan,
84 resolved_ids: ResolvedIds,
85 ) {
86 let stage = return_if_err!(
87 self.create_index_validate(ctx.session(), plan, resolved_ids, ExplainContext::None),
88 ctx
89 );
90 self.sequence_staged(ctx, Span::current(), stage).await;
91 }
92
93 #[instrument]
94 pub(crate) async fn explain_create_index(
95 &mut self,
96 ctx: ExecuteContext,
97 plan::ExplainPlanPlan {
98 stage,
99 format,
100 config,
101 explainee,
102 }: plan::ExplainPlanPlan,
103 ) {
104 let plan::Explainee::Statement(stmt) = explainee else {
105 unreachable!()
108 };
109 let plan::ExplaineeStatement::CreateIndex { broken, plan } = stmt else {
110 unreachable!()
113 };
114
115 let optimizer_trace = OptimizerTrace::new(stage.paths());
118
119 let resolved_ids = ResolvedIds::empty();
121
122 let explain_ctx = ExplainContext::Plan(ExplainPlanContext {
123 broken,
124 config,
125 format,
126 stage,
127 replan: None,
128 desc: None,
129 optimizer_trace,
130 });
131 let stage = return_if_err!(
132 self.create_index_validate(ctx.session(), plan, resolved_ids, explain_ctx),
133 ctx
134 );
135 self.sequence_staged(ctx, Span::current(), stage).await;
136 }
137
138 #[instrument]
139 pub(crate) async fn explain_replan_index(
140 &mut self,
141 ctx: ExecuteContext,
142 plan::ExplainPlanPlan {
143 stage,
144 format,
145 config,
146 explainee,
147 }: plan::ExplainPlanPlan,
148 ) {
149 let plan::Explainee::ReplanIndex(id) = explainee else {
150 unreachable!() };
152 let CatalogItem::Index(index) = self.catalog().get_entry(&id).item() else {
153 unreachable!() };
155 let id = index.global_id();
156
157 let create_sql = index.create_sql.clone();
158 let plan_result = self
159 .catalog_mut()
160 .deserialize_plan_with_enable_for_item_parsing(&create_sql, true);
161 let (plan, resolved_ids) = return_if_err!(plan_result, ctx);
162
163 let plan::Plan::CreateIndex(plan) = plan else {
164 unreachable!() };
166
167 let broken = false;
170
171 let optimizer_trace = OptimizerTrace::new(stage.paths());
174
175 let explain_ctx = ExplainContext::Plan(ExplainPlanContext {
176 broken,
177 config,
178 format,
179 stage,
180 replan: Some(id),
181 desc: None,
182 optimizer_trace,
183 });
184 let stage = return_if_err!(
185 self.create_index_validate(ctx.session(), plan, resolved_ids, explain_ctx),
186 ctx
187 );
188 self.sequence_staged(ctx, Span::current(), stage).await;
189 }
190
191 #[instrument]
192 pub(crate) fn explain_index(
193 &self,
194 ctx: &ExecuteContext,
195 plan::ExplainPlanPlan {
196 stage,
197 format,
198 config,
199 explainee,
200 }: plan::ExplainPlanPlan,
201 ) -> Result<ExecuteResponse, AdapterError> {
202 let plan::Explainee::Index(id) = explainee else {
203 unreachable!() };
205 let CatalogItem::Index(index) = self.catalog().get_entry(&id).item() else {
206 unreachable!() };
208
209 let Some(dataflow_metainfo) = self.catalog().try_get_dataflow_metainfo(&index.global_id())
210 else {
211 if !id.is_system() {
212 tracing::error!("cannot find dataflow metainformation for index {id} in catalog");
213 }
214 coord_bail!("cannot find dataflow metainformation for index {id} in catalog");
215 };
216
217 let target_cluster = self.catalog().get_cluster(index.cluster_id);
218
219 let features = OptimizerFeatures::from(self.catalog().system_config())
220 .override_from(&target_cluster.config.features())
221 .override_from(&self.cluster_scoped_optimizer_overrides(index.cluster_id))
222 .override_from(&config.features);
223
224 let cardinality_stats = BTreeMap::new();
226
227 let explain = match stage {
228 ExplainStage::GlobalPlan => {
229 let Some(plan) = self
230 .catalog()
231 .try_get_optimized_plan(&index.global_id())
232 .cloned()
233 else {
234 tracing::error!("cannot find {stage} for index {id} in catalog");
235 coord_bail!("cannot find {stage} for index in catalog");
236 };
237
238 explain_dataflow(
239 plan,
240 format,
241 &config,
242 &features,
243 &self.catalog().for_session(ctx.session()),
244 cardinality_stats,
245 Some(target_cluster.name.as_str()),
246 dataflow_metainfo,
247 )?
248 }
249 ExplainStage::PhysicalPlan => {
250 let Some(plan) = self
251 .catalog()
252 .try_get_physical_plan(&index.global_id())
253 .cloned()
254 else {
255 tracing::error!("cannot find {stage} for index {id} in catalog");
256 coord_bail!("cannot find {stage} for index in catalog");
257 };
258 explain_dataflow(
259 plan,
260 format,
261 &config,
262 &features,
263 &self.catalog().for_session(ctx.session()),
264 cardinality_stats,
265 Some(target_cluster.name.as_str()),
266 dataflow_metainfo,
267 )?
268 }
269 _ => {
270 coord_bail!("cannot EXPLAIN {} FOR INDEX", stage);
271 }
272 };
273
274 let row = Row::pack_slice(&[Datum::from(explain.as_str())]);
275
276 Ok(Self::send_immediate_rows(row))
277 }
278
279 #[instrument]
282 fn create_index_validate(
283 &self,
284 session: &Session,
285 plan: plan::CreateIndexPlan,
286 resolved_ids: ResolvedIds,
287 explain_ctx: ExplainContext,
288 ) -> Result<CreateIndexStage, AdapterError> {
289 let validity = PlanValidity::new(
293 self.catalog(),
294 resolved_ids.items().copied().collect(),
295 Some(plan.index.cluster_id),
296 None,
297 session.role_metadata().clone(),
298 );
299 Ok(CreateIndexStage::Optimize(CreateIndexOptimize {
300 validity,
301 plan,
302 resolved_ids,
303 explain_ctx,
304 }))
305 }
306
307 #[instrument]
308 async fn create_index_optimize(
309 &mut self,
310 CreateIndexOptimize {
311 validity,
312 plan,
313 resolved_ids,
314 explain_ctx,
315 }: CreateIndexOptimize,
316 ) -> Result<StageResult<Box<CreateIndexStage>>, AdapterError> {
317 let plan::CreateIndexPlan {
318 index: plan::Index { cluster_id, .. },
319 ..
320 } = &plan;
321
322 let compute_instance = self
324 .instance_snapshot(*cluster_id)
325 .expect("compute instance does not exist");
326 let (item_id, global_id) = if let ExplainContext::None = explain_ctx {
327 self.allocate_user_id().await?
328 } else {
329 self.allocate_transient_id()
330 };
331
332 let optimizer_config = optimize::OptimizerConfig::from(self.catalog().system_config())
333 .override_from(&self.catalog.get_cluster(*cluster_id).config.features())
334 .override_from(&self.cluster_scoped_optimizer_overrides(*cluster_id))
335 .override_from(&explain_ctx);
336 let optimizer_features = optimizer_config.features.clone();
337
338 let mut optimizer = optimize::index::Optimizer::new(
340 self.owned_catalog(),
341 compute_instance,
342 global_id,
343 optimizer_config,
344 self.optimizer_metrics(),
345 );
346 let span = Span::current();
347 Ok(StageResult::Handle(mz_ore::task::spawn_blocking(
348 || "optimize create index",
349 move || {
350 span.in_scope(|| {
351 let mut pipeline = || -> Result<(
352 optimize::index::GlobalMirPlan,
353 optimize::index::GlobalLirPlan,
354 ), AdapterError> {
355 let _dispatch_guard = explain_ctx.dispatch_guard();
356
357 let index_plan = optimize::index::Index::new(
358 plan.name.clone(),
359 plan.index.on,
360 plan.index.keys.clone(),
361 );
362
363 let global_mir_plan = optimizer.catch_unwind_optimize(index_plan)?;
365 let global_lir_plan = optimizer.catch_unwind_optimize(global_mir_plan.clone())?;
367
368 Ok((global_mir_plan, global_lir_plan))
369 };
370
371 let stage = match pipeline() {
372 Ok((global_mir_plan, global_lir_plan)) => {
373 if let ExplainContext::Plan(explain_ctx) = explain_ctx {
374 let (_, df_meta) = global_lir_plan.unapply();
375 CreateIndexStage::Explain(CreateIndexExplain {
376 validity,
377 exported_index_id: global_id,
378 plan,
379 df_meta,
380 explain_ctx,
381 })
382 } else {
383 CreateIndexStage::Finish(CreateIndexFinish {
384 validity,
385 item_id,
386 global_id,
387 plan,
388 resolved_ids,
389 global_mir_plan,
390 global_lir_plan,
391 optimizer_features,
392 })
393 }
394 }
395 Err(err) => {
398 let ExplainContext::Plan(explain_ctx) = explain_ctx else {
399 return Err(err);
401 };
402
403 if explain_ctx.broken {
404 tracing::error!("error while handling EXPLAIN statement: {}", err);
408 CreateIndexStage::Explain(CreateIndexExplain {
409 validity,
410 exported_index_id: global_id,
411 plan,
412 df_meta: Default::default(),
413 explain_ctx,
414 })
415 } else {
416 return Err(err);
418 }
419 }
420 };
421 Ok(Box::new(stage))
422 })
423 },
424 )))
425 }
426
427 #[instrument]
428 async fn create_index_finish(
429 &mut self,
430 ctx: &mut ExecuteContext,
431 stage: CreateIndexFinish,
432 ) -> Result<StageResult<Box<CreateIndexStage>>, AdapterError> {
433 let CreateIndexFinish {
434 item_id,
435 global_id,
436 plan:
437 plan::CreateIndexPlan {
438 name,
439 index:
440 plan::Index {
441 create_sql,
442 on,
443 keys,
444 cluster_id,
445 compaction_window,
446 },
447 if_not_exists,
448 },
449 resolved_ids,
450 global_mir_plan,
451 global_lir_plan,
452 optimizer_features,
453 ..
454 } = stage;
455 let id_bundle = dataflow_import_id_bundle(global_lir_plan.df_desc(), cluster_id);
456
457 let on_entry = self.catalog().get_entry_by_global_id(&on);
458 let owner_id = *on_entry.owner_id();
459
460 let ops = vec![catalog::Op::CreateItem {
461 id: item_id,
462 name: name.clone(),
463 item: CatalogItem::Index(Index {
464 create_sql,
465 global_id,
466 keys: keys.into(),
467 on,
468 conn_id: None,
469 resolved_ids,
470 cluster_id,
471 is_retained_metrics_object: false,
472 custom_logical_compaction_window: compaction_window,
473 optimized_plan: None,
474 physical_plan: None,
475 dataflow_metainfo: None,
476 }),
477 owner_id,
478 }];
479
480 let (df_desc, raw_df_meta) = global_lir_plan.unapply();
484 let on_desc = on_entry
485 .relation_desc()
486 .expect("can only create indexes on items with a valid description");
487 let df_meta = self.render_create_item_notices(&name, global_id, &on_desc, &raw_df_meta);
488
489 self.catalog()
494 .cache_expressions(
495 global_id,
496 None,
497 global_mir_plan.df_desc().clone(),
498 df_desc.clone(),
499 df_meta.clone(),
500 optimizer_features,
501 )
502 .await;
503
504 let transact_result = self
505 .catalog_transact_with_side_effects(Some(ctx), ops, move |coord, _ctx| {
506 Box::pin(async move {
507 coord
509 .catalog_mut()
510 .set_optimized_plan(global_id, global_mir_plan.df_desc().clone());
511 coord
512 .catalog_mut()
513 .set_physical_plan(global_id, df_desc.clone());
514
515 let notice_builtin_updates_fut =
516 coord.persist_dataflow_metainfo(df_meta, global_id);
517
518 coord
522 .ship_new_dataflow(
523 &id_bundle,
524 df_desc,
525 cluster_id,
526 notice_builtin_updates_fut,
527 )
528 .await;
529 coord.update_compute_read_policy(
532 cluster_id,
533 item_id,
534 compaction_window.unwrap_or_default().into(),
535 );
536 })
537 })
538 .await;
539
540 match transact_result {
541 Ok(_) => {
542 self.emit_raw_optimizer_notices_to_user(ctx, &raw_df_meta.optimizer_notices);
547 Ok(StageResult::Response(ExecuteResponse::CreatedIndex))
548 }
549 Err(AdapterError::Catalog(mz_catalog::memory::error::Error {
550 kind: ErrorKind::Sql(CatalogError::ItemAlreadyExists(_, _)),
551 })) if if_not_exists => {
552 ctx.session()
553 .add_notice(AdapterNotice::ObjectAlreadyExists {
554 name: name.item,
555 ty: "index",
556 });
557 Ok(StageResult::Response(ExecuteResponse::CreatedIndex))
558 }
559 Err(err) => Err(err),
560 }
561 }
562
563 #[instrument]
564 async fn create_index_explain(
565 &self,
566 session: &Session,
567 CreateIndexExplain {
568 exported_index_id,
569 plan: plan::CreateIndexPlan { name, index, .. },
570 df_meta,
571 explain_ctx:
572 ExplainPlanContext {
573 config,
574 format,
575 stage,
576 optimizer_trace,
577 ..
578 },
579 ..
580 }: CreateIndexExplain,
581 ) -> Result<StageResult<Box<CreateIndexStage>>, AdapterError> {
582 let session_catalog = self.catalog().for_session(session);
583 let expr_humanizer = {
584 let on_entry = self.catalog.get_entry_by_global_id(&index.on);
585 let full_name = self.catalog.resolve_full_name(&name, on_entry.conn_id());
586 let on_desc = on_entry
587 .relation_desc()
588 .expect("can only create indexes on items with a valid description");
589
590 let transient_items = btreemap! {
591 exported_index_id => TransientItem::new(
592 Some(full_name.into_parts()),
593 Some(on_desc.iter_names().map(|c| c.to_string()).collect()),
594 )
595 };
596 ExprHumanizerExt::new(transient_items, &session_catalog)
597 };
598
599 let target_cluster = self.catalog().get_cluster(index.cluster_id);
600
601 let features = OptimizerFeatures::from(self.catalog().system_config())
602 .override_from(&target_cluster.config.features())
603 .override_from(&self.cluster_scoped_optimizer_overrides(index.cluster_id))
604 .override_from(&config.features);
605
606 let rows = optimizer_trace
607 .into_rows(
608 format,
609 &config,
610 &features,
611 &expr_humanizer,
612 None,
613 Some(target_cluster),
614 df_meta,
615 stage,
616 plan::ExplaineeStatementKind::CreateIndex,
617 None,
618 )
619 .await?;
620
621 Ok(StageResult::Response(Self::send_immediate_rows(rows)))
622 }
623}