Skip to main content

mz_adapter/coord/sequencer/inner/
create_index.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10use 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            // This is currently asserted in the `sequence_explain_plan` code that
106            // calls this method.
107            unreachable!()
108        };
109        let plan::ExplaineeStatement::CreateIndex { broken, plan } = stmt else {
110            // This is currently asserted in the `sequence_explain_plan` code that
111            // calls this method.
112            unreachable!()
113        };
114
115        // Create an OptimizerTrace instance to collect plans emitted when
116        // executing the optimizer pipeline.
117        let optimizer_trace = OptimizerTrace::new(stage.paths());
118
119        // Not used in the EXPLAIN path so it's OK to generate a dummy value.
120        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!() // Asserted in `sequence_explain_plan`.
151        };
152        let CatalogItem::Index(index) = self.catalog().get_entry(&id).item() else {
153            unreachable!() // Asserted in `plan_explain_plan`.
154        };
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!() // We are parsing the `create_sql` of an `Index` item.
165        };
166
167        // It is safe to assume that query optimization will always succeed, so
168        // for now we statically assume `broken = false`.
169        let broken = false;
170
171        // Create an OptimizerTrace instance to collect plans emitted when
172        // executing the optimizer pipeline.
173        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!() // Asserted in `sequence_explain_plan`.
204        };
205        let CatalogItem::Index(index) = self.catalog().get_entry(&id).item() else {
206            unreachable!() // Asserted in `plan_explain_plan`.
207        };
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        // TODO(mgree): calculate statistics (need a timestamp)
225        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    // `explain_ctx` is an optional context set iff the state machine is initiated from
280    // sequencing an EXPLAIN for this statement.
281    #[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        // Track the target cluster and resolved dependencies so concurrent
290        // drops are caught between stages instead of panicking later when the
291        // persisted SQL is re-parsed during catalog application.
292        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        // Collect optimizer parameters.
323        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        // Build an optimizer for this INDEX.
339        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                    // MIR ⇒ MIR optimization (global)
364                    let global_mir_plan = optimizer.catch_unwind_optimize(index_plan)?;
365                    // MIR ⇒ LIR lowering and LIR ⇒ LIR optimization (global)
366                    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                        // Internal optimizer errors are handled differently
396                        // depending on the caller.
397                        Err(err) => {
398                            let ExplainContext::Plan(explain_ctx) = explain_ctx else {
399                                // In `sequence_~` contexts, immediately error.
400                                return Err(err);
401                            };
402
403                            if explain_ctx.broken {
404                                // In `EXPLAIN BROKEN` contexts, just log the error
405                                // and move to the next stage with default
406                                // parameters.
407                                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                                // In regular `EXPLAIN` contexts, immediately error.
417                                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        // We keep `raw_df_meta` live so that on success we can emit its raw
481        // notices to the user session (rendered against the user's
482        // session-aware humanizer).
483        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        // Populate the durable expression cache before the catalog
490        // transaction and await the write. This way any other envd (or a
491        // subsequent bootstrap here) will observe the cached plans +
492        // rendered notices as soon as the item becomes visible.
493        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                    // Save plan structures.
508                    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                    // TODO: Maybe in the future, pass the read holds
519                    // `ship_new_dataflow` takes on to compute, to hold on to them
520                    // and downgrade when possible?
521                    coord
522                        .ship_new_dataflow(
523                            &id_bundle,
524                            df_desc,
525                            cluster_id,
526                            notice_builtin_updates_fut,
527                        )
528                        .await;
529                    // No `allow_writes` here because indexes do not modify external state.
530
531                    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                // Only emit optimizer notices to the user now that the
543                // catalog transaction has succeeded. If the transaction had
544                // failed, emitting notices would confuse the user with
545                // information about an item that wasn't actually created.
546                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}