Skip to main content

mz_adapter/
coord.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
10//! Translation of SQL commands into timestamped `Controller` commands.
11//!
12//! The various SQL commands instruct the system to take actions that are not
13//! yet explicitly timestamped. On the other hand, the underlying data continually
14//! change as time moves forward. On the third hand, we greatly benefit from the
15//! information that some times are no longer of interest, so that we may
16//! compact the representation of the continually changing collections.
17//!
18//! The [`Coordinator`] curates these interactions by observing the progress
19//! collections make through time, choosing timestamps for its own commands,
20//! and eventually communicating that certain times have irretrievably "passed".
21//!
22//! ## Frontiers another way
23//!
24//! If the above description of frontiers left you with questions, this
25//! repackaged explanation might help.
26//!
27//! - `since` is the least recent time (i.e. oldest time) that you can read
28//!   from sources and be guaranteed that the returned data is accurate as of
29//!   that time.
30//!
31//!   Reads at times less than `since` may return values that were not actually
32//!   seen at the specified time, but arrived later (i.e. the results are
33//!   compacted).
34//!
35//!   For correctness' sake, the coordinator never chooses to read at a time
36//!   less than an arrangement's `since`.
37//!
38//! - `upper` is the first time after the most recent time that you can read
39//!   from sources and receive an immediate response. Alternately, it is the
40//!   least time at which the data may still change (that is the reason we may
41//!   not be able to respond immediately).
42//!
43//!   Reads at times >= `upper` may not immediately return because the answer
44//!   isn't known yet. However, once the `upper` is > the specified read time,
45//!   the read can return.
46//!
47//!   For the sake of returned values' freshness, the coordinator prefers
48//!   performing reads at an arrangement's `upper`. However, because we more
49//!   strongly prefer correctness, the coordinator will choose timestamps
50//!   greater than an object's `upper` if it is also being accessed alongside
51//!   objects whose `since` times are >= its `upper`.
52//!
53//! This illustration attempts to show, with time moving left to right, the
54//! relationship between `since` and `upper`.
55//!
56//! - `#`: possibly inaccurate results
57//! - `-`: immediate, correct response
58//! - `?`: not yet known
59//! - `s`: since
60//! - `u`: upper
61//! - `|`: eligible for coordinator to select
62//!
63//! ```nofmt
64//! ####s----u?????
65//!     |||||||||||
66//! ```
67//!
68
69use std::borrow::Cow;
70use std::collections::{BTreeMap, BTreeSet, VecDeque};
71use std::net::IpAddr;
72use std::num::NonZeroI64;
73use std::ops::Neg;
74use std::str::FromStr;
75use std::sync::LazyLock;
76use std::sync::{Arc, Mutex};
77use std::thread;
78use std::time::{Duration, Instant};
79use std::{fmt, mem};
80
81use anyhow::Context;
82use chrono::{DateTime, Utc};
83use derivative::Derivative;
84use differential_dataflow::lattice::Lattice;
85use fail::fail_point;
86use futures::StreamExt;
87use futures::future::{BoxFuture, FutureExt, LocalBoxFuture};
88use http::Uri;
89use ipnet::IpNet;
90use itertools::Itertools;
91use mz_adapter_types::bootstrap_builtin_cluster_config::BootstrapBuiltinClusterConfig;
92use mz_adapter_types::compaction::CompactionWindow;
93use mz_adapter_types::connection::ConnectionId;
94use mz_adapter_types::dyncfgs::FRONTEND_READ_THEN_WRITE;
95use mz_adapter_types::dyncfgs::{
96    ENABLE_0DT_HYDRATE_MIGRATED_BUILTIN_MVS, USER_ID_POOL_BATCH_SIZE,
97    WITH_0DT_DEPLOYMENT_CAUGHT_UP_CHECK_INTERVAL,
98};
99use mz_auth::password::Password;
100use mz_build_info::BuildInfo;
101use mz_catalog::builtin::{
102    BUILTINS, BUILTINS_STATIC, MZ_OBJECT_ARRANGEMENT_SIZE_HISTORY, MZ_OBJECT_HYDRATION_HISTORY,
103    MZ_REPLICA_HYDRATION_HISTORY, MZ_STORAGE_USAGE_BY_SHARD,
104};
105use mz_catalog::config::{AwsPrincipalContext, BuiltinItemMigrationConfig, ClusterReplicaSizeMap};
106use mz_catalog::durable::OpenableDurableCatalogState;
107use mz_catalog::expr_cache::{GlobalExpressions, LocalExpressions, latest_item_version};
108use mz_catalog::memory::objects::{
109    CatalogEntry, CatalogItem, ClusterReplicaProcessStatus, Connection, DataSourceDesc,
110    ReconfigurationTarget, Table, TableDataSource,
111};
112use mz_cloud_resources::{CloudResourceController, VpcEndpointConfig, VpcEndpointEvent};
113use mz_compute_client::as_of_selection;
114use mz_compute_client::controller::error::{
115    CollectionLookupError, CollectionMissing, DataflowCreationError, InstanceMissing,
116};
117use mz_compute_types::ComputeInstanceId;
118use mz_compute_types::dataflows::DataflowDescription;
119use mz_compute_types::plan::LirRelationExpr;
120use mz_controller::clusters::{
121    ClusterConfig, ClusterEvent, ClusterStatus, ManagedReplicaLocation, ProcessId, ReplicaLocation,
122};
123use mz_controller::{ControllerConfig, Readiness};
124use mz_controller_types::{ClusterId, ReplicaId, WatchSetId};
125use mz_dyncfg::{ConfigUpdates, ParameterScope};
126use mz_expr::{MapFilterProject, MirRelationExpr, OptimizedMirRelationExpr, RowSetFinishing};
127use mz_license_keys::{ExpirationBehavior, ValidatedLicenseKey};
128use mz_orchestrator::OfflineReason;
129use mz_ore::cast::{CastFrom, CastInto, CastLossy};
130use mz_ore::channel::trigger::Trigger;
131use mz_ore::future::TimeoutError;
132use mz_ore::metrics::MetricsRegistry;
133use mz_ore::now::{EpochMillis, NowFn};
134use mz_ore::task::{AbortOnDropHandle, JoinHandle, spawn};
135use mz_ore::thread::JoinHandleExt;
136use mz_ore::tracing::{OpenTelemetryContext, TracingHandle};
137use mz_ore::url::SensitiveUrl;
138use mz_ore::{
139    assert_none, instrument, soft_assert_eq_or_log, soft_assert_or_log, soft_panic_or_log, stack,
140};
141use mz_persist_client::PersistClient;
142use mz_persist_client::batch::ProtoBatch;
143use mz_persist_client::usage::{ShardsUsageReferenced, StorageUsageClient};
144use mz_repr::adt::numeric::Numeric;
145use mz_repr::explain::{ExplainConfig, ExplainFormat};
146use mz_repr::global_id::TransientIdGen;
147use mz_repr::optimize::{OptimizerFeatureOverrides, OptimizerFeatures, OverrideFrom};
148use mz_repr::role_id::RoleId;
149use mz_repr::{
150    CatalogItemId, Diff, GlobalId, RelationDesc, RelationVersion, SqlRelationType, Timestamp,
151};
152use mz_secrets::cache::CachingSecretsReader;
153use mz_secrets::{SecretsController, SecretsReader};
154use mz_sql::ast::{Raw, Statement};
155use mz_sql::catalog::{CatalogCluster, EnvironmentId};
156use mz_sql::names::{QualifiedItemName, ResolvedIds};
157use mz_sql::optimizer_metrics::OptimizerMetrics;
158use mz_sql::plan::{
159    self, AlterSinkPlan, ConnectionDetails, CreateConnectionPlan, HirRelationExpr,
160    NetworkPolicyRule, Params, QueryWhen,
161};
162use mz_sql::session::user::User;
163use mz_sql::session::vars::{MAX_CREDIT_CONSUMPTION_RATE, SystemVars, Var};
164use mz_sql_parser::ast::ExplainStage;
165use mz_sql_parser::ast::display::AstDisplay;
166use mz_storage_client::client::TableData;
167use mz_storage_client::controller::{CollectionDescription, DataSource, ExportDescription};
168use mz_storage_types::connections::Connection as StorageConnection;
169use mz_storage_types::connections::ConnectionContext;
170use mz_storage_types::connections::inline::{IntoInlineConnection, ReferencedConnection};
171use mz_storage_types::read_holds::ReadHold;
172use mz_storage_types::sinks::{S3SinkFormat, StorageSinkDesc};
173use mz_storage_types::sources::kafka::KAFKA_PROGRESS_DESC;
174use mz_storage_types::sources::{IngestionDescription, SourceExport, Timeline};
175use mz_timestamp_oracle::{TimestampOracleConfig, WriteTimestamp};
176use mz_transform::dataflow::DataflowMetainfo;
177use opentelemetry::trace::TraceContextExt;
178use semver::Version;
179use serde::Serialize;
180use thiserror::Error;
181use timely::progress::{Antichain, Timestamp as _};
182use tokio::runtime::Handle as TokioHandle;
183use tokio::select;
184use tokio::sync::{Notify, OwnedMutexGuard, Semaphore, mpsc, oneshot, watch};
185use tokio::time::{Interval, MissedTickBehavior};
186use tracing::{Instrument, Level, Span, debug, info, info_span, span, warn};
187use tracing_opentelemetry::OpenTelemetrySpanExt;
188use uuid::Uuid;
189
190use crate::active_compute_sink::{ActiveComputeSink, ActiveCopyFrom};
191use crate::catalog::{BuiltinTableUpdate, Catalog, OpenCatalogResult};
192use crate::client::{Client, Handle};
193use crate::command::{Command, ExecuteResponse};
194use crate::config::{
195    ClusterEvalContext, ClusterScopeContext, ReplicaEvalContext, ReplicaScopeContext,
196    ScopedParameters, ScopedParametersScope, SynchronizedParameters, SystemParameterFrontend,
197    SystemParameterSyncConfig,
198};
199use crate::coord::appends::{
200    BuiltinTableAppendCompletion, BuiltinTableAppendNotify, DeferredOp, GroupCommitPermit,
201    PendingWriteTxn,
202};
203use crate::coord::caught_up::CaughtUpCheckContext;
204use crate::coord::id_bundle::CollectionIdBundle;
205use crate::coord::introspection::IntrospectionSubscribe;
206use crate::coord::metric_sink::{CuratedMetricSink, InstalledMetricSink, PlannedMetricSink};
207use crate::coord::peek::PendingPeek;
208use crate::coord::statement_logging::StatementLogging;
209use crate::coord::timeline::{TimelineContext, TimelineState};
210use crate::coord::timestamp_selection::{TimestampContext, TimestampDetermination};
211use crate::coord::validity::PlanValidity;
212use crate::error::AdapterError;
213use crate::explain::insights::PlanInsightsContext;
214use crate::explain::optimizer_trace::{DispatchGuard, OptimizerTrace};
215use crate::metrics::Metrics;
216use crate::optimize::dataflows::{ComputeInstanceSnapshot, DataflowBuilder};
217use crate::optimize::{self, Optimize, OptimizerConfig};
218use crate::session::{EndTransactionAction, Session};
219use crate::statement_logging::{
220    StatementEndedExecutionReason, StatementLifecycleEvent, StatementLoggingId,
221};
222use crate::util::{ClientTransmitter, ResultExt, sort_topological};
223use crate::webhook::{WebhookAppenderInvalidator, WebhookConcurrencyLimiter};
224use crate::{AdapterNotice, ReadHolds, flags};
225
226pub(crate) mod appends;
227pub(crate) mod catalog_serving;
228pub(crate) mod cluster_controller;
229pub(crate) mod consistency;
230pub(crate) mod id_bundle;
231pub(crate) mod in_memory_oracle;
232pub(crate) mod peek;
233pub(crate) mod read_policy;
234pub(crate) mod read_then_write;
235pub(crate) mod sequencer;
236pub(crate) mod statement_logging;
237pub(crate) mod timeline;
238pub(crate) mod timestamp_selection;
239
240pub mod catalog_implications;
241mod caught_up;
242mod command_handler;
243mod ddl;
244pub(crate) mod group_sync;
245mod hydration_history;
246mod indexes;
247mod info_metrics;
248mod introspection;
249mod message_handler;
250mod metric_sink;
251mod privatelink_status;
252mod sql;
253mod validity;
254
255/// The oldest leader version against which a replacement-migrated builtin materialized view may
256/// write its new persist shard while this environment is still read-only.
257///
258/// Every builtin materialized view reads `mz_internal.mz_catalog_raw`, so its dataflow only makes
259/// progress up to the catalog shard's frontier. Holding that frontier at the current time is the
260/// leader's job, and leaders only started doing it in v26.17 (PR #35402). Write-enable such an MV
261/// against an older leader and it sits at a stale frontier and never reports caught up, which
262/// blocks promotion outright instead of merely leaving the collection cold at cut-over. We still
263/// support upgrading from before v26.17, so that leader is a real case, not a hypothetical.
264const MIN_LEADER_VERSION_FOR_MIGRATED_MV_WRITES: Version = Version::new(26, 17, 0);
265
266/// A pool of pre-allocated user IDs to avoid per-DDL persist writes.
267///
268/// IDs in the range `[next, upper)` are available for allocation.
269/// When exhausted, the pool must be refilled via the catalog.
270///
271/// # Correctness
272///
273/// The pool is owned by [`Coordinator`], which processes all requests
274/// on a single-threaded event loop. Because every access requires
275/// `&mut self` on the coordinator, there is no concurrent access to the
276/// pool — no additional synchronization is needed.
277///
278/// Global ID uniqueness is guaranteed because each refill calls
279/// [`Catalog::allocate_user_ids`], which performs a durable persist
280/// write that atomically reserves the entire batch before any IDs from
281/// it are handed out. If the process crashes after a refill but before
282/// all pre-allocated IDs are consumed, the unused IDs form harmless
283/// gaps in the sequence — user IDs are not required to be contiguous.
284///
285/// This guarantee holds even if multiple `environmentd` processes run
286/// concurrently. Each process has its own independent pool,
287/// but every refill goes through the shared persist-backed catalog,
288/// which serializes allocations across all callers. Two processes
289/// will therefore never receive overlapping ID ranges,
290/// for the same reason they could not before this pool existed.
291#[derive(Debug)]
292pub(crate) struct IdPool {
293    next: u64,
294    upper: u64,
295}
296
297impl IdPool {
298    /// Creates an empty pool.
299    pub fn empty() -> Self {
300        IdPool { next: 0, upper: 0 }
301    }
302
303    /// Allocates a single ID from the pool, returning `None` if exhausted.
304    pub fn allocate(&mut self) -> Option<u64> {
305        if self.next < self.upper {
306            let id = self.next;
307            self.next += 1;
308            Some(id)
309        } else {
310            None
311        }
312    }
313
314    /// Allocates `n` consecutive IDs from the pool, returning `None` if
315    /// insufficient IDs remain.
316    pub fn allocate_many(&mut self, n: u64) -> Option<Vec<u64>> {
317        if self.remaining() >= n {
318            let ids = (self.next..self.next + n).collect();
319            self.next += n;
320            Some(ids)
321        } else {
322            None
323        }
324    }
325
326    /// Returns the number of IDs remaining in the pool.
327    pub fn remaining(&self) -> u64 {
328        self.upper - self.next
329    }
330
331    /// Refills the pool with the given range `[next, upper)`.
332    pub fn refill(&mut self, next: u64, upper: u64) {
333        assert!(next <= upper, "invalid pool range: {next}..{upper}");
334        self.next = next;
335        self.upper = upper;
336    }
337}
338
339/// A row for `mz_object_arrangement_size_history`, prepared off-thread by the
340/// arrangement sizes snapshot task and stamped with a collection timestamp at
341/// write time.
342#[derive(Debug)]
343pub struct ArrangementSizeRecord {
344    pub replica_id: String,
345    pub object_id: String,
346    pub size: i64,
347    pub hydration_complete: bool,
348}
349
350#[derive(Debug)]
351pub enum Message {
352    Command(OpenTelemetryContext, Command),
353    ControllerReady {
354        controller: ControllerReadiness,
355    },
356    PurifiedStatementReady(PurifiedStatementReady),
357    CreateConnectionValidationReady(CreateConnectionValidationReady),
358    AlterConnectionValidationReady(AlterConnectionValidationReady),
359    TryDeferred {
360        /// The connection that created this op.
361        conn_id: ConnectionId,
362        /// The write lock that notified us our deferred op might be able to run.
363        ///
364        /// Note: While we never want to hold a partial set of locks, it can be important to hold
365        /// onto the _one_ that notified us our op might be ready. If there are multiple operations
366        /// waiting on a single collection, and we don't hold this lock through retyring the op,
367        /// then everything waiting on this collection will get retried causing traffic in the
368        /// Coordinator's message queue.
369        ///
370        /// See [`DeferredOp::can_be_optimistically_retried`] for more detail.
371        acquired_lock: Option<(CatalogItemId, tokio::sync::OwnedMutexGuard<()>)>,
372    },
373    /// Initiates a group commit.
374    GroupCommitInitiate(Span, Option<GroupCommitPermit>),
375    /// Finalizes an applied group commit.
376    ///
377    /// Statement timestamps precede response retirement because retirement ends statement logging.
378    GroupCommitApplied {
379        /// Responses to retire after recording statement timestamps.
380        responses: Vec<crate::util::CompletedClientTransmitter>,
381        /// Statement executions associated with this commit.
382        statement_logging_ids: Vec<StatementLoggingId>,
383        /// Frontend-sequenced writes to complete after local timestamp bookkeeping.
384        internal_results: Vec<crate::coord::appends::InternalWriteResponder>,
385        /// The applied write timestamp.
386        write_ts: Timestamp,
387    },
388    DeferredStatementReady,
389    AdvanceTimelines,
390    ClusterEvent(ClusterEvent),
391    CancelPendingPeeks {
392        conn_id: ConnectionId,
393    },
394    LinearizeReads,
395    StagedBatches {
396        conn_id: ConnectionId,
397        table_id: CatalogItemId,
398        batches: Vec<Result<ProtoBatch, String>>,
399    },
400    StorageUsageSchedule,
401    StorageUsageFetch,
402    StorageUsageUpdate(ShardsUsageReferenced),
403    StorageUsagePrune(Vec<BuiltinTableUpdate>),
404    ArrangementSizesSchedule,
405    ArrangementSizesSnapshot,
406    ArrangementSizesWrite(Vec<ArrangementSizeRecord>),
407    ArrangementSizesPrune(Vec<BuiltinTableUpdate>),
408    HydrationHistorySchedule,
409    HydrationHistoryRun,
410    /// Performs any cleanup and logging actions necessary for
411    /// finalizing a statement execution.
412    RetireExecute {
413        data: ExecuteContextExtra,
414        otel_ctx: OpenTelemetryContext,
415        reason: StatementEndedExecutionReason,
416    },
417    ExecuteSingleStatementTransaction {
418        ctx: ExecuteContext,
419        otel_ctx: OpenTelemetryContext,
420        stmt: Arc<Statement<Raw>>,
421        params: mz_sql::plan::Params,
422    },
423    PeekStageReady {
424        ctx: ExecuteContext,
425        span: Span,
426        stage: PeekStage,
427    },
428    CreateIndexStageReady {
429        ctx: ExecuteContext,
430        span: Span,
431        stage: CreateIndexStage,
432    },
433    CreateMetricSinkStageReady {
434        ctx: ExecuteContext,
435        span: Span,
436        stage: CreateMetricSinkStage,
437    },
438    CreateViewStageReady {
439        ctx: ExecuteContext,
440        span: Span,
441        stage: CreateViewStage,
442    },
443    CreateMaterializedViewStageReady {
444        ctx: ExecuteContext,
445        span: Span,
446        stage: CreateMaterializedViewStage,
447    },
448    SubscribeStageReady {
449        ctx: ExecuteContext,
450        span: Span,
451        stage: SubscribeStage,
452    },
453    IntrospectionSubscribeStageReady {
454        span: Span,
455        stage: IntrospectionSubscribeStage,
456    },
457    MetricSinkStageReady {
458        span: Span,
459        stage: MetricSinkStage,
460    },
461    SecretStageReady {
462        ctx: ExecuteContext,
463        span: Span,
464        stage: SecretStage,
465    },
466    ClusterStageReady {
467        ctx: ExecuteContext,
468        span: Span,
469        stage: ClusterStage,
470    },
471    ExplainTimestampStageReady {
472        ctx: ExecuteContext,
473        span: Span,
474        stage: ExplainTimestampStage,
475    },
476    DrainStatementLog,
477    PrivateLinkVpcEndpointEvents(Vec<VpcEndpointEvent>),
478
479    /// One pull/apply call from the cluster controller task, answered on the main
480    /// coordinator message loop from the catalog and live controller signals.
481    /// See [`cluster_controller`].
482    ClusterControllerRequest(cluster_controller::ClusterControllerRequest),
483}
484
485impl Message {
486    /// Returns a string to identify the kind of [`Message`], useful for logging.
487    pub const fn kind(&self) -> &'static str {
488        match self {
489            Message::Command(_, msg) => match msg {
490                Command::CatalogSnapshot { .. } => "command-catalog_snapshot",
491                Command::Startup { .. } => "command-startup",
492                Command::Execute { .. } => "command-execute",
493                Command::Commit { .. } => "command-commit",
494                Command::CancelRequest { .. } => "command-cancel_request",
495                Command::PrivilegedCancelRequest { .. } => "command-privileged_cancel_request",
496                Command::GetWebhook { .. } => "command-get_webhook",
497                Command::GetSystemVars { .. } => "command-get_system_vars",
498                Command::SetSystemVars { .. } => "command-set_system_vars",
499                Command::UpdateScopedSystemParameters { .. } => {
500                    "command-update_scoped_system_parameters"
501                }
502                Command::InstallScopedSystemParameterFrontend { .. } => {
503                    "command-install_scoped_system_parameter_frontend"
504                }
505                Command::Terminate { .. } => "command-terminate",
506                Command::RetireExecute { .. } => "command-retire_execute",
507                Command::CheckConsistency { .. } => "command-check_consistency",
508                Command::Dump { .. } => "command-dump",
509                Command::AuthenticatePassword { .. } => "command-auth_check",
510                Command::AuthenticateGetSASLChallenge { .. } => "command-auth_get_sasl_challenge",
511                Command::AuthenticateVerifySASLProof { .. } => "command-auth_verify_sasl_proof",
512                Command::CheckRoleCanLogin { .. } => "command-check_role_can_login",
513                Command::GetComputeInstanceClient { .. } => "get-compute-instance-client",
514                Command::GetOracle { .. } => "get-oracle",
515                Command::DetermineRealTimeRecentTimestamp { .. } => {
516                    "determine-real-time-recent-timestamp"
517                }
518                Command::GetTransactionReadHoldsBundle { .. } => {
519                    "get-transaction-read-holds-bundle"
520                }
521                Command::StoreTransactionReadHolds { .. } => "store-transaction-read-holds",
522                Command::ExecuteSlowPathPeek { .. } => "execute-slow-path-peek",
523                Command::ExecuteSubscribe { .. } => "execute-subscribe",
524                Command::CopyToPreflight { .. } => "copy-to-preflight",
525                Command::ExecuteCopyTo { .. } => "execute-copy-to",
526                Command::ExecuteSideEffectingFunc { .. } => "execute-side-effecting-func",
527                Command::LookupConnection { .. } => "lookup-connection",
528                Command::RegisterFrontendPeek { .. } => "register-frontend-peek",
529                Command::UnregisterFrontendPeek { .. } => "unregister-frontend-peek",
530                Command::ExplainTimestamp { .. } => "explain-timestamp",
531                Command::FrontendStatementLogging(..) => "frontend-statement-logging",
532                Command::StartCopyFromStdin { .. } => "start-copy-from-stdin",
533                Command::InjectAuditEvents { .. } => "inject-audit-events",
534                Command::RegisterConnectionCancelWatch { .. } => "register-connection-cancel-watch",
535                Command::CreateInternalSubscribe { .. } => "create-internal-subscribe",
536                Command::AttemptWrite { .. } => "attempt-write",
537                Command::DropInternalSubscribe { .. } => "drop-internal-subscribe",
538            },
539            Message::ControllerReady {
540                controller: ControllerReadiness::Compute,
541            } => "controller_ready(compute)",
542            Message::ControllerReady {
543                controller: ControllerReadiness::Storage,
544            } => "controller_ready(storage)",
545            Message::ControllerReady {
546                controller: ControllerReadiness::Metrics,
547            } => "controller_ready(metrics)",
548            Message::ControllerReady {
549                controller: ControllerReadiness::Internal,
550            } => "controller_ready(internal)",
551            Message::PurifiedStatementReady(_) => "purified_statement_ready",
552            Message::CreateConnectionValidationReady(_) => "create_connection_validation_ready",
553            Message::TryDeferred { .. } => "try_deferred",
554            Message::GroupCommitInitiate(..) => "group_commit_initiate",
555            Message::GroupCommitApplied { .. } => "group_commit_applied",
556            Message::AdvanceTimelines => "advance_timelines",
557            Message::ClusterEvent(_) => "cluster_event",
558            Message::CancelPendingPeeks { .. } => "cancel_pending_peeks",
559            Message::LinearizeReads => "linearize_reads",
560            Message::StagedBatches { .. } => "staged_batches",
561            Message::StorageUsageSchedule => "storage_usage_schedule",
562            Message::StorageUsageFetch => "storage_usage_fetch",
563            Message::StorageUsageUpdate(_) => "storage_usage_update",
564            Message::StorageUsagePrune(_) => "storage_usage_prune",
565            Message::ArrangementSizesSchedule => "arrangement_sizes_schedule",
566            Message::ArrangementSizesSnapshot => "arrangement_sizes_snapshot",
567            Message::ArrangementSizesWrite(_) => "arrangement_sizes_write",
568            Message::ArrangementSizesPrune(_) => "arrangement_sizes_prune",
569            Message::HydrationHistorySchedule => "hydration_history_schedule",
570            Message::HydrationHistoryRun => "hydration_history_run",
571            Message::RetireExecute { .. } => "retire_execute",
572            Message::ExecuteSingleStatementTransaction { .. } => {
573                "execute_single_statement_transaction"
574            }
575            Message::PeekStageReady { .. } => "peek_stage_ready",
576            Message::ExplainTimestampStageReady { .. } => "explain_timestamp_stage_ready",
577            Message::CreateIndexStageReady { .. } => "create_index_stage_ready",
578            Message::CreateMetricSinkStageReady { .. } => "create_metric_sink_stage_ready",
579            Message::CreateViewStageReady { .. } => "create_view_stage_ready",
580            Message::CreateMaterializedViewStageReady { .. } => {
581                "create_materialized_view_stage_ready"
582            }
583            Message::SubscribeStageReady { .. } => "subscribe_stage_ready",
584            Message::IntrospectionSubscribeStageReady { .. } => {
585                "introspection_subscribe_stage_ready"
586            }
587            Message::MetricSinkStageReady { .. } => "metric_sink_stage_ready",
588            Message::SecretStageReady { .. } => "secret_stage_ready",
589            Message::ClusterStageReady { .. } => "cluster_stage_ready",
590            Message::DrainStatementLog => "drain_statement_log",
591            Message::AlterConnectionValidationReady(..) => "alter_connection_validation_ready",
592            Message::PrivateLinkVpcEndpointEvents(_) => "private_link_vpc_endpoint_events",
593            Message::ClusterControllerRequest(_) => "cluster_controller_request",
594            Message::DeferredStatementReady => "deferred_statement_ready",
595        }
596    }
597}
598
599/// The reason for why a controller needs processing on the main loop.
600#[derive(Debug)]
601pub enum ControllerReadiness {
602    /// The storage controller is ready.
603    Storage,
604    /// The compute controller is ready.
605    Compute,
606    /// A batch of metric data is ready.
607    Metrics,
608    /// An internally-generated message is ready to be returned.
609    Internal,
610}
611
612#[derive(Derivative)]
613#[derivative(Debug)]
614pub struct BackgroundWorkResult<T> {
615    #[derivative(Debug = "ignore")]
616    pub ctx: ExecuteContext,
617    pub result: Result<T, AdapterError>,
618    pub params: Params,
619    pub plan_validity: PlanValidity,
620    pub original_stmt: Arc<Statement<Raw>>,
621    pub otel_ctx: OpenTelemetryContext,
622}
623
624pub type PurifiedStatementReady = BackgroundWorkResult<mz_sql::pure::PurifiedStatement>;
625
626#[derive(Derivative)]
627#[derivative(Debug)]
628pub struct ValidationReady<T> {
629    #[derivative(Debug = "ignore")]
630    pub ctx: ExecuteContext,
631    pub result: Result<T, AdapterError>,
632    pub resolved_ids: ResolvedIds,
633    pub connection_id: CatalogItemId,
634    pub connection_gid: GlobalId,
635    pub plan_validity: PlanValidity,
636    pub otel_ctx: OpenTelemetryContext,
637}
638
639pub type CreateConnectionValidationReady = ValidationReady<CreateConnectionPlan>;
640pub type AlterConnectionValidationReady = ValidationReady<Connection>;
641
642#[derive(Debug)]
643pub enum PeekStage {
644    /// Common stages across SELECT, EXPLAIN and COPY TO queries.
645    LinearizeTimestamp(PeekStageLinearizeTimestamp),
646    RealTimeRecency(PeekStageRealTimeRecency),
647    TimestampReadHold(PeekStageTimestampReadHold),
648    Optimize(PeekStageOptimize),
649    /// Final stage for a peek.
650    Finish(PeekStageFinish),
651    /// Final stage for an explain.
652    ExplainPlan(PeekStageExplainPlan),
653    ExplainPushdown(PeekStageExplainPushdown),
654    /// Preflight checks for a copy to operation.
655    CopyToPreflight(PeekStageCopyTo),
656    /// Final stage for a copy to which involves shipping the dataflow.
657    CopyToDataflow(PeekStageCopyTo),
658}
659
660#[derive(Debug)]
661pub struct CopyToContext {
662    /// The `RelationDesc` of the data to be copied.
663    pub desc: RelationDesc,
664    /// The destination uri of the external service where the data will be copied.
665    pub uri: Uri,
666    /// Connection information required to connect to the external service to copy the data.
667    pub connection: StorageConnection<ReferencedConnection>,
668    /// The ID of the CONNECTION object to be used for copying the data.
669    pub connection_id: CatalogItemId,
670    /// Format params to format the data.
671    pub format: S3SinkFormat,
672    /// Approximate max file size of each uploaded file.
673    pub max_file_size: u64,
674    /// Number of batches the output of the COPY TO will be partitioned into
675    /// to distribute the load across workers deterministically.
676    /// This is only an option since it's not set when CopyToContext is instantiated
677    /// but immediately after in the PeekStageValidate stage.
678    pub output_batch_count: Option<u64>,
679}
680
681#[derive(Debug)]
682pub struct PeekStageLinearizeTimestamp {
683    validity: PlanValidity,
684    plan: mz_sql::plan::SelectPlan,
685    max_query_result_size: Option<u64>,
686    source_ids: BTreeSet<GlobalId>,
687    target_replica: Option<ReplicaId>,
688    timeline_context: TimelineContext,
689    optimizer: optimize::PeekOptimizer,
690    /// An optional context set iff the state machine is initiated from
691    /// sequencing an EXPLAIN for this statement.
692    explain_ctx: ExplainContext,
693}
694
695#[derive(Debug)]
696pub struct PeekStageRealTimeRecency {
697    validity: PlanValidity,
698    plan: mz_sql::plan::SelectPlan,
699    max_query_result_size: Option<u64>,
700    source_ids: BTreeSet<GlobalId>,
701    target_replica: Option<ReplicaId>,
702    timeline_context: TimelineContext,
703    oracle_read_ts: Option<Timestamp>,
704    optimizer: optimize::PeekOptimizer,
705    /// An optional context set iff the state machine is initiated from
706    /// sequencing an EXPLAIN for this statement.
707    explain_ctx: ExplainContext,
708}
709
710#[derive(Debug)]
711pub struct PeekStageTimestampReadHold {
712    validity: PlanValidity,
713    plan: mz_sql::plan::SelectPlan,
714    max_query_result_size: Option<u64>,
715    source_ids: BTreeSet<GlobalId>,
716    target_replica: Option<ReplicaId>,
717    timeline_context: TimelineContext,
718    oracle_read_ts: Option<Timestamp>,
719    real_time_recency_ts: Option<mz_repr::Timestamp>,
720    optimizer: optimize::PeekOptimizer,
721    /// An optional context set iff the state machine is initiated from
722    /// sequencing an EXPLAIN for this statement.
723    explain_ctx: ExplainContext,
724}
725
726#[derive(Debug)]
727pub struct PeekStageOptimize {
728    validity: PlanValidity,
729    plan: mz_sql::plan::SelectPlan,
730    max_query_result_size: Option<u64>,
731    source_ids: BTreeSet<GlobalId>,
732    id_bundle: CollectionIdBundle,
733    target_replica: Option<ReplicaId>,
734    determination: TimestampDetermination,
735    optimizer: optimize::PeekOptimizer,
736    /// An optional context set iff the state machine is initiated from
737    /// sequencing an EXPLAIN for this statement.
738    explain_ctx: ExplainContext,
739}
740
741#[derive(Debug)]
742pub struct PeekStageFinish {
743    validity: PlanValidity,
744    plan: mz_sql::plan::SelectPlan,
745    max_query_result_size: Option<u64>,
746    id_bundle: CollectionIdBundle,
747    target_replica: Option<ReplicaId>,
748    source_ids: BTreeSet<GlobalId>,
749    determination: TimestampDetermination,
750    cluster_id: ComputeInstanceId,
751    finishing: RowSetFinishing,
752    /// When present, an optimizer trace to be used for emitting a plan insights
753    /// notice.
754    plan_insights_optimizer_trace: Option<OptimizerTrace>,
755    insights_ctx: Option<Box<PlanInsightsContext>>,
756    global_lir_plan: optimize::peek::GlobalLirPlan,
757    optimization_finished_at: EpochMillis,
758}
759
760#[derive(Debug)]
761pub struct PeekStageCopyTo {
762    validity: PlanValidity,
763    optimizer: optimize::copy_to::Optimizer,
764    global_lir_plan: optimize::copy_to::GlobalLirPlan,
765    optimization_finished_at: EpochMillis,
766    target_replica: Option<ReplicaId>,
767    source_ids: BTreeSet<GlobalId>,
768}
769
770#[derive(Debug)]
771pub struct PeekStageExplainPlan {
772    validity: PlanValidity,
773    optimizer: optimize::peek::Optimizer,
774    df_meta: DataflowMetainfo,
775    explain_ctx: ExplainPlanContext,
776    insights_ctx: Option<Box<PlanInsightsContext>>,
777}
778
779#[derive(Debug)]
780pub struct PeekStageExplainPushdown {
781    validity: PlanValidity,
782    determination: TimestampDetermination,
783    imports: BTreeMap<GlobalId, MapFilterProject>,
784}
785
786#[derive(Debug)]
787pub enum CreateIndexStage {
788    Optimize(CreateIndexOptimize),
789    Finish(CreateIndexFinish),
790    Explain(CreateIndexExplain),
791}
792
793#[derive(Debug)]
794pub struct CreateIndexOptimize {
795    validity: PlanValidity,
796    plan: plan::CreateIndexPlan,
797    resolved_ids: ResolvedIds,
798    /// An optional context set iff the state machine is initiated from
799    /// sequencing an EXPLAIN for this statement.
800    explain_ctx: ExplainContext,
801}
802
803#[derive(Debug)]
804pub struct CreateIndexFinish {
805    validity: PlanValidity,
806    item_id: CatalogItemId,
807    global_id: GlobalId,
808    plan: plan::CreateIndexPlan,
809    resolved_ids: ResolvedIds,
810    global_mir_plan: optimize::index::GlobalMirPlan,
811    global_lir_plan: optimize::index::GlobalLirPlan,
812    optimizer_features: OptimizerFeatures,
813}
814
815#[derive(Debug)]
816pub struct CreateIndexExplain {
817    validity: PlanValidity,
818    exported_index_id: GlobalId,
819    plan: plan::CreateIndexPlan,
820    df_meta: DataflowMetainfo,
821    explain_ctx: ExplainPlanContext,
822}
823
824#[derive(Debug)]
825pub enum CreateMetricSinkStage {
826    Optimize(CreateMetricSinkOptimize),
827    Finish(CreateMetricSinkFinish),
828}
829
830#[derive(Debug)]
831pub struct CreateMetricSinkOptimize {
832    validity: PlanValidity,
833    plan: plan::CreateMetricSinkPlan,
834    resolved_ids: ResolvedIds,
835}
836
837#[derive(Debug)]
838pub struct CreateMetricSinkFinish {
839    validity: PlanValidity,
840    item_id: CatalogItemId,
841    global_id: GlobalId,
842    plan: plan::CreateMetricSinkPlan,
843    resolved_ids: ResolvedIds,
844    global_mir_plan: optimize::metric_sink::GlobalMirPlan,
845    global_lir_plan: optimize::metric_sink::GlobalLirPlan,
846    optimizer_features: OptimizerFeatures,
847}
848
849#[derive(Debug)]
850pub enum CreateViewStage {
851    Optimize(CreateViewOptimize),
852    Finish(CreateViewFinish),
853    Explain(CreateViewExplain),
854}
855
856#[derive(Debug)]
857pub struct CreateViewOptimize {
858    validity: PlanValidity,
859    plan: plan::CreateViewPlan,
860    resolved_ids: ResolvedIds,
861    /// An optional context set iff the state machine is initiated from
862    /// sequencing an EXPLAIN for this statement.
863    explain_ctx: ExplainContext,
864}
865
866#[derive(Debug)]
867pub struct CreateViewFinish {
868    validity: PlanValidity,
869    /// ID of this item in the Catalog.
870    item_id: CatalogItemId,
871    /// ID by with Compute will reference this View.
872    global_id: GlobalId,
873    plan: plan::CreateViewPlan,
874    /// IDs of objects resolved during name resolution.
875    resolved_ids: ResolvedIds,
876    optimized_expr: OptimizedMirRelationExpr,
877}
878
879#[derive(Debug)]
880pub struct CreateViewExplain {
881    validity: PlanValidity,
882    id: GlobalId,
883    plan: plan::CreateViewPlan,
884    explain_ctx: ExplainPlanContext,
885}
886
887#[derive(Debug)]
888pub enum ExplainTimestampStage {
889    Optimize(ExplainTimestampOptimize),
890    RealTimeRecency(ExplainTimestampRealTimeRecency),
891    LinearizeTimestamp(ExplainTimestampLinearizeTimestamp),
892    Finish(ExplainTimestampFinish),
893}
894
895#[derive(Debug)]
896pub struct ExplainTimestampOptimize {
897    validity: PlanValidity,
898    plan: plan::ExplainTimestampPlan,
899    cluster_id: ClusterId,
900}
901
902#[derive(Debug)]
903pub struct ExplainTimestampRealTimeRecency {
904    validity: PlanValidity,
905    format: ExplainFormat,
906    optimized_plan: OptimizedMirRelationExpr,
907    cluster_id: ClusterId,
908    when: QueryWhen,
909}
910
911#[derive(Debug)]
912pub struct ExplainTimestampLinearizeTimestamp {
913    validity: PlanValidity,
914    format: ExplainFormat,
915    optimized_plan: OptimizedMirRelationExpr,
916    cluster_id: ClusterId,
917    source_ids: BTreeSet<GlobalId>,
918    when: QueryWhen,
919    real_time_recency_ts: Option<Timestamp>,
920}
921
922#[derive(Debug)]
923pub struct ExplainTimestampFinish {
924    validity: PlanValidity,
925    format: ExplainFormat,
926    cluster_id: ClusterId,
927    source_ids: BTreeSet<GlobalId>,
928    when: QueryWhen,
929    real_time_recency_ts: Option<Timestamp>,
930    /// The timeline context derived in the preceding `LinearizeTimestamp`
931    /// stage, carried forward so it stays consistent with `oracle_read_ts`.
932    timeline_context: TimelineContext,
933    /// The linearized read timestamp, read off the coordinator loop in the
934    /// preceding `LinearizeTimestamp` stage. `None` when no linearized read is
935    /// needed.
936    oracle_read_ts: Option<Timestamp>,
937}
938
939#[derive(Debug)]
940pub enum ClusterStage {
941    Alter(AlterCluster),
942    /// The foreground wait-shim over a controller-driven background
943    /// reconfiguration: poll the durable `reconfiguration` record until it
944    /// clears, then report success or timeout depending on whether the realized
945    /// config reached the target.
946    AwaitReconfiguration(AlterClusterAwaitReconfiguration),
947}
948
949#[derive(Debug)]
950pub struct AlterCluster {
951    validity: PlanValidity,
952    plan: plan::AlterClusterPlan,
953}
954
955#[derive(Debug)]
956pub struct AlterClusterAwaitReconfiguration {
957    validity: PlanValidity,
958    cluster_id: ClusterId,
959    /// The target shape the awaited `ALTER` wrote. Once the record becomes
960    /// terminal, the realized config matching this is what distinguishes a
961    /// cut-over from a failure. See `await_reconfiguration_stage`.
962    target: ReconfigurationTarget,
963}
964
965#[derive(Debug)]
966pub enum ExplainContext {
967    /// The ordinary, non-explain variant of the statement.
968    None,
969    /// The `EXPLAIN <level> PLAN FOR <explainee>` version of the statement.
970    Plan(ExplainPlanContext),
971    /// Generate a notice containing the `EXPLAIN PLAN INSIGHTS` output
972    /// alongside the query's normal output.
973    PlanInsightsNotice(OptimizerTrace),
974    /// `EXPLAIN FILTER PUSHDOWN`
975    Pushdown,
976}
977
978impl ExplainContext {
979    /// If available for this context, wrap the [`OptimizerTrace`] into a
980    /// [`tracing::Dispatch`] and set it as default, returning the resulting
981    /// guard in a `Some(guard)` option.
982    pub(crate) fn dispatch_guard(&self) -> Option<DispatchGuard<'_>> {
983        let optimizer_trace = match self {
984            ExplainContext::Plan(explain_ctx) => Some(&explain_ctx.optimizer_trace),
985            ExplainContext::PlanInsightsNotice(optimizer_trace) => Some(optimizer_trace),
986            _ => None,
987        };
988        optimizer_trace.map(|optimizer_trace| optimizer_trace.as_guard())
989    }
990
991    pub(crate) fn needs_cluster(&self) -> bool {
992        match self {
993            ExplainContext::None => true,
994            ExplainContext::Plan(..) => false,
995            ExplainContext::PlanInsightsNotice(..) => true,
996            ExplainContext::Pushdown => false,
997        }
998    }
999
1000    pub(crate) fn needs_plan_insights(&self) -> bool {
1001        matches!(
1002            self,
1003            ExplainContext::Plan(ExplainPlanContext {
1004                stage: ExplainStage::PlanInsights,
1005                ..
1006            }) | ExplainContext::PlanInsightsNotice(_)
1007        )
1008    }
1009}
1010
1011#[derive(Debug)]
1012pub struct ExplainPlanContext {
1013    /// EXPLAIN BROKEN is internal syntax for showing EXPLAIN output despite an internal error in
1014    /// the optimizer: we don't immediately bail out from peek sequencing when an internal optimizer
1015    /// error happens, but go on with trying to show the requested EXPLAIN stage. This can still
1016    /// succeed if the requested EXPLAIN stage is before the point where the error happened.
1017    pub broken: bool,
1018    pub config: ExplainConfig,
1019    pub format: ExplainFormat,
1020    pub stage: ExplainStage,
1021    pub replan: Option<GlobalId>,
1022    pub desc: Option<RelationDesc>,
1023    pub optimizer_trace: OptimizerTrace,
1024}
1025
1026#[derive(Debug)]
1027pub enum CreateMaterializedViewStage {
1028    Optimize(CreateMaterializedViewOptimize),
1029    Finish(CreateMaterializedViewFinish),
1030    Explain(CreateMaterializedViewExplain),
1031}
1032
1033#[derive(Debug)]
1034pub struct CreateMaterializedViewOptimize {
1035    validity: PlanValidity,
1036    plan: plan::CreateMaterializedViewPlan,
1037    resolved_ids: ResolvedIds,
1038    /// An optional context set iff the state machine is initiated from
1039    /// sequencing an EXPLAIN for this statement.
1040    explain_ctx: ExplainContext,
1041}
1042
1043#[derive(Debug)]
1044pub struct CreateMaterializedViewFinish {
1045    /// The ID of this Materialized View in the Catalog.
1046    item_id: CatalogItemId,
1047    /// The ID of the durable pTVC backing this Materialized View.
1048    global_id: GlobalId,
1049    validity: PlanValidity,
1050    plan: plan::CreateMaterializedViewPlan,
1051    resolved_ids: ResolvedIds,
1052    local_mir_plan: optimize::materialized_view::LocalMirPlan,
1053    global_mir_plan: optimize::materialized_view::GlobalMirPlan,
1054    global_lir_plan: optimize::materialized_view::GlobalLirPlan,
1055    optimizer_features: OptimizerFeatures,
1056}
1057
1058#[derive(Debug)]
1059pub struct CreateMaterializedViewExplain {
1060    global_id: GlobalId,
1061    validity: PlanValidity,
1062    plan: plan::CreateMaterializedViewPlan,
1063    df_meta: DataflowMetainfo,
1064    explain_ctx: ExplainPlanContext,
1065}
1066
1067#[derive(Debug)]
1068pub enum SubscribeStage {
1069    OptimizeMir(SubscribeOptimizeMir),
1070    LinearizeTimestamp(SubscribeLinearizeTimestamp),
1071    TimestampOptimizeLir(SubscribeTimestampOptimizeLir),
1072    Finish(SubscribeFinish),
1073    Explain(SubscribeExplain),
1074}
1075
1076#[derive(Debug)]
1077pub struct SubscribeOptimizeMir {
1078    validity: PlanValidity,
1079    plan: plan::SubscribePlan,
1080    timeline: TimelineContext,
1081    dependency_ids: BTreeSet<GlobalId>,
1082    cluster_id: ComputeInstanceId,
1083    replica_id: Option<ReplicaId>,
1084    /// An optional context set iff the state machine is initiated from
1085    /// sequencing an EXPLAIN for this statement.
1086    explain_ctx: ExplainContext,
1087}
1088
1089#[derive(Debug)]
1090pub struct SubscribeLinearizeTimestamp {
1091    validity: PlanValidity,
1092    plan: plan::SubscribePlan,
1093    timeline: TimelineContext,
1094    optimizer: optimize::subscribe::Optimizer,
1095    global_mir_plan: optimize::subscribe::GlobalMirPlan<optimize::subscribe::Unresolved>,
1096    dependency_ids: BTreeSet<GlobalId>,
1097    replica_id: Option<ReplicaId>,
1098    /// An optional context set iff the state machine is initiated from
1099    /// sequencing an EXPLAIN for this statement.
1100    explain_ctx: ExplainContext,
1101}
1102
1103#[derive(Debug)]
1104pub struct SubscribeTimestampOptimizeLir {
1105    validity: PlanValidity,
1106    plan: plan::SubscribePlan,
1107    timeline: TimelineContext,
1108    optimizer: optimize::subscribe::Optimizer,
1109    global_mir_plan: optimize::subscribe::GlobalMirPlan<optimize::subscribe::Unresolved>,
1110    dependency_ids: BTreeSet<GlobalId>,
1111    replica_id: Option<ReplicaId>,
1112    /// The linearized read timestamp, read off the coordinator loop in the
1113    /// preceding `LinearizeTimestamp` stage. `None` when no linearized read is
1114    /// needed.
1115    oracle_read_ts: Option<Timestamp>,
1116    /// An optional context set iff the state machine is initiated from
1117    /// sequencing an EXPLAIN for this statement.
1118    explain_ctx: ExplainContext,
1119}
1120
1121#[derive(Debug)]
1122pub struct SubscribeFinish {
1123    validity: PlanValidity,
1124    cluster_id: ComputeInstanceId,
1125    replica_id: Option<ReplicaId>,
1126    plan: plan::SubscribePlan,
1127    global_lir_plan: optimize::subscribe::GlobalLirPlan,
1128    dependency_ids: BTreeSet<GlobalId>,
1129}
1130
1131#[derive(Debug)]
1132pub struct SubscribeExplain {
1133    validity: PlanValidity,
1134    optimizer: optimize::subscribe::Optimizer,
1135    df_meta: DataflowMetainfo,
1136    cluster_id: ComputeInstanceId,
1137    explain_ctx: ExplainPlanContext,
1138}
1139
1140#[derive(Debug)]
1141pub enum IntrospectionSubscribeStage {
1142    OptimizeMir(IntrospectionSubscribeOptimizeMir),
1143    TimestampOptimizeLir(IntrospectionSubscribeTimestampOptimizeLir),
1144    Finish(IntrospectionSubscribeFinish),
1145}
1146
1147#[derive(Debug)]
1148pub struct IntrospectionSubscribeOptimizeMir {
1149    validity: PlanValidity,
1150    plan: plan::SubscribePlan,
1151    subscribe_id: GlobalId,
1152    cluster_id: ComputeInstanceId,
1153    replica_id: ReplicaId,
1154}
1155
1156#[derive(Debug)]
1157pub struct IntrospectionSubscribeTimestampOptimizeLir {
1158    validity: PlanValidity,
1159    optimizer: optimize::subscribe::Optimizer,
1160    global_mir_plan: optimize::subscribe::GlobalMirPlan<optimize::subscribe::Unresolved>,
1161    cluster_id: ComputeInstanceId,
1162    replica_id: ReplicaId,
1163}
1164
1165#[derive(Debug)]
1166pub struct IntrospectionSubscribeFinish {
1167    validity: PlanValidity,
1168    global_lir_plan: optimize::subscribe::GlobalLirPlan,
1169    read_holds: ReadHolds,
1170    cluster_id: ComputeInstanceId,
1171    replica_id: ReplicaId,
1172}
1173
1174#[derive(Debug)]
1175pub enum MetricSinkStage {
1176    Optimize(MetricSinkOptimize),
1177    Finish(MetricSinkFinish),
1178}
1179
1180#[derive(Debug)]
1181pub struct MetricSinkOptimize {
1182    validity: PlanValidity,
1183    definition: &'static CuratedMetricSink,
1184    /// The transient id of the sink's compute export. Recorded in
1185    /// [`Coordinator::metric_sinks`] once the finish stage ships the dataflow.
1186    sink_id: GlobalId,
1187    /// The planned `source_sql`, and the shape it produces.
1188    expr: HirRelationExpr,
1189    desc: RelationDesc,
1190    cluster_id: ComputeInstanceId,
1191    replica_id: ReplicaId,
1192}
1193
1194#[derive(Debug)]
1195pub struct MetricSinkFinish {
1196    validity: PlanValidity,
1197    definition: &'static CuratedMetricSink,
1198    sink_id: GlobalId,
1199    global_lir_plan: optimize::metric_sink::GlobalLirPlan,
1200    cluster_id: ComputeInstanceId,
1201    replica_id: ReplicaId,
1202}
1203
1204#[derive(Debug)]
1205pub enum SecretStage {
1206    CreateEnsure(CreateSecretEnsure),
1207    CreateFinish(CreateSecretFinish),
1208    RotateKeysEnsure(RotateKeysSecretEnsure),
1209    RotateKeysFinish(RotateKeysSecretFinish),
1210    Alter(AlterSecret),
1211}
1212
1213#[derive(Debug)]
1214pub struct CreateSecretEnsure {
1215    validity: PlanValidity,
1216    plan: plan::CreateSecretPlan,
1217}
1218
1219#[derive(Debug)]
1220pub struct CreateSecretFinish {
1221    validity: PlanValidity,
1222    item_id: CatalogItemId,
1223    global_id: GlobalId,
1224    plan: plan::CreateSecretPlan,
1225}
1226
1227#[derive(Debug)]
1228pub struct RotateKeysSecretEnsure {
1229    validity: PlanValidity,
1230    id: CatalogItemId,
1231}
1232
1233#[derive(Debug)]
1234pub struct RotateKeysSecretFinish {
1235    validity: PlanValidity,
1236    ops: Vec<crate::catalog::Op>,
1237}
1238
1239#[derive(Debug)]
1240pub struct AlterSecret {
1241    validity: PlanValidity,
1242    plan: plan::AlterSecretPlan,
1243}
1244
1245/// An enum describing which cluster to run a statement on.
1246///
1247/// One example usage would be that if a query depends only on system tables, we might
1248/// automatically run it on the catalog server cluster to benefit from indexes that exist there.
1249#[derive(Debug, Copy, Clone, PartialEq, Eq)]
1250pub enum TargetCluster {
1251    /// The catalog server cluster.
1252    CatalogServer,
1253    /// The current user's active cluster.
1254    Active,
1255    /// The cluster selected at the start of a transaction.
1256    Transaction(ClusterId),
1257}
1258
1259/// Result types for each stage of a sequence.
1260pub(crate) enum StageResult<T> {
1261    /// A task was spawned that will return the next stage.
1262    Handle(JoinHandle<Result<T, AdapterError>>),
1263    /// A task was spawned that will return a response for the client.
1264    HandleRetire(JoinHandle<Result<ExecuteResponse, AdapterError>>),
1265    /// The next stage is immediately ready and will execute.
1266    Immediate(T),
1267    /// The final stage was executed and is ready to respond to the client.
1268    Response(ExecuteResponse),
1269}
1270
1271/// Common functionality for [Coordinator::sequence_staged].
1272pub(crate) trait Staged: Send {
1273    type Ctx: StagedContext;
1274
1275    fn validity(&mut self) -> &mut PlanValidity;
1276
1277    /// Returns the next stage or final result.
1278    async fn stage(
1279        self,
1280        coord: &mut Coordinator,
1281        ctx: &mut Self::Ctx,
1282    ) -> Result<StageResult<Box<Self>>, AdapterError>;
1283
1284    /// Prepares a message for the Coordinator.
1285    fn message(self, ctx: Self::Ctx, span: Span) -> Message;
1286
1287    /// Whether it is safe to SQL cancel this stage.
1288    fn cancel_enabled(&self) -> bool;
1289}
1290
1291pub trait StagedContext {
1292    fn retire(self, result: Result<ExecuteResponse, AdapterError>);
1293    fn session(&self) -> Option<&Session>;
1294}
1295
1296impl StagedContext for ExecuteContext {
1297    fn retire(self, result: Result<ExecuteResponse, AdapterError>) {
1298        self.retire(result);
1299    }
1300
1301    fn session(&self) -> Option<&Session> {
1302        Some(self.session())
1303    }
1304}
1305
1306impl StagedContext for () {
1307    fn retire(self, _result: Result<ExecuteResponse, AdapterError>) {}
1308
1309    fn session(&self) -> Option<&Session> {
1310        None
1311    }
1312}
1313
1314/// Configures a coordinator.
1315pub struct Config {
1316    pub controller_config: ControllerConfig,
1317    pub controller_envd_epoch: NonZeroI64,
1318    pub storage: Box<dyn mz_catalog::durable::DurableCatalogState>,
1319    pub timestamp_oracle_url: Option<SensitiveUrl>,
1320    pub unsafe_mode: bool,
1321    pub all_features: bool,
1322    pub build_info: &'static BuildInfo,
1323    pub environment_id: EnvironmentId,
1324    pub metrics_registry: MetricsRegistry,
1325    pub now: NowFn,
1326    pub secrets_controller: Arc<dyn SecretsController>,
1327    pub cloud_resource_controller: Option<Arc<dyn CloudResourceController>>,
1328    pub availability_zones: Vec<String>,
1329    pub cluster_replica_sizes: ClusterReplicaSizeMap,
1330    pub builtin_system_cluster_config: BootstrapBuiltinClusterConfig,
1331    pub builtin_catalog_server_cluster_config: BootstrapBuiltinClusterConfig,
1332    pub builtin_probe_cluster_config: BootstrapBuiltinClusterConfig,
1333    pub builtin_support_cluster_config: BootstrapBuiltinClusterConfig,
1334    pub builtin_analytics_cluster_config: BootstrapBuiltinClusterConfig,
1335    pub system_parameter_defaults: BTreeMap<String, String>,
1336    pub storage_usage_client: StorageUsageClient,
1337    pub storage_usage_collection_interval: Duration,
1338    pub storage_usage_retention_period: Option<Duration>,
1339    pub segment_client: Option<mz_segment::Client>,
1340    pub egress_addresses: Vec<IpNet>,
1341    pub remote_system_parameters: Option<BTreeMap<String, String>>,
1342    pub aws_account_id: Option<String>,
1343    pub aws_privatelink_availability_zones: Option<Vec<String>>,
1344    pub connection_context: ConnectionContext,
1345    pub connection_limit_callback: Box<dyn Fn(u64, u64) -> () + Send + Sync + 'static>,
1346    pub webhook_concurrency_limit: WebhookConcurrencyLimiter,
1347    pub http_host_name: Option<String>,
1348    pub tracing_handle: TracingHandle,
1349    /// Whether or not to start controllers in read-only mode. This is only
1350    /// meant for use during development of read-only clusters and 0dt upgrades
1351    /// and should go away once we have proper orchestration during upgrades.
1352    pub read_only_controllers: bool,
1353
1354    /// A trigger that signals that the current deployment has caught up with a
1355    /// previous deployment. Only used during 0dt deployment, while in read-only
1356    /// mode.
1357    pub caught_up_trigger: Option<Trigger>,
1358
1359    pub helm_chart_version: Option<String>,
1360    pub license_key: ValidatedLicenseKey,
1361    pub external_login_password_mz_system: Option<Password>,
1362    pub force_builtin_schema_migration: Option<String>,
1363}
1364
1365/// Metadata about an active connection.
1366#[derive(Debug, Serialize)]
1367pub struct ConnMeta {
1368    /// Pgwire specifies that every connection have a 32-bit secret associated
1369    /// with it, that is known to both the client and the server. Cancellation
1370    /// requests are required to authenticate with the secret of the connection
1371    /// that they are targeting.
1372    secret_key: u32,
1373    /// The time when the session's connection was initiated.
1374    connected_at: EpochMillis,
1375    user: User,
1376    application_name: String,
1377    uuid: Uuid,
1378    conn_id: ConnectionId,
1379    client_ip: Option<IpAddr>,
1380
1381    /// Sinks that will need to be dropped when the current transaction, if
1382    /// any, is cleared.
1383    drop_sinks: BTreeSet<GlobalId>,
1384
1385    /// Lock for the Coordinator's deferred statements that is dropped on transaction clear.
1386    #[serde(skip)]
1387    deferred_lock: Option<OwnedMutexGuard<()>>,
1388
1389    /// Channel on which to send notices to a session.
1390    #[serde(skip)]
1391    notice_tx: mpsc::UnboundedSender<AdapterNotice>,
1392
1393    /// The role that initiated the database context. Fixed for the duration of the connection.
1394    /// WARNING: This role reference is not updated when the role is dropped.
1395    /// Consumers should not assume that this role exist.
1396    authenticated_role: RoleId,
1397}
1398
1399impl ConnMeta {
1400    pub fn conn_id(&self) -> &ConnectionId {
1401        &self.conn_id
1402    }
1403
1404    pub fn user(&self) -> &User {
1405        &self.user
1406    }
1407
1408    pub fn application_name(&self) -> &str {
1409        &self.application_name
1410    }
1411
1412    pub fn authenticated_role_id(&self) -> &RoleId {
1413        &self.authenticated_role
1414    }
1415
1416    pub fn uuid(&self) -> Uuid {
1417        self.uuid
1418    }
1419
1420    pub fn client_ip(&self) -> Option<IpAddr> {
1421        self.client_ip
1422    }
1423
1424    pub fn connected_at(&self) -> EpochMillis {
1425        self.connected_at
1426    }
1427}
1428
1429#[derive(Debug)]
1430/// A pending transaction waiting to be committed.
1431pub struct PendingTxn {
1432    /// Context used to send a response back to the client.
1433    ctx: ExecuteContext,
1434    /// Client response for transaction.
1435    response: Result<PendingTxnResponse, AdapterError>,
1436    /// The action to take at the end of the transaction.
1437    action: EndTransactionAction,
1438}
1439
1440#[derive(Debug)]
1441/// The response we'll send for a [`PendingTxn`].
1442pub enum PendingTxnResponse {
1443    /// The transaction will be committed.
1444    Committed {
1445        /// Parameters that will change, and their values, once this transaction is complete.
1446        params: BTreeMap<&'static str, String>,
1447    },
1448    /// The transaction will be rolled back.
1449    Rolledback {
1450        /// Parameters that will change, and their values, once this transaction is complete.
1451        params: BTreeMap<&'static str, String>,
1452    },
1453}
1454
1455impl PendingTxnResponse {
1456    pub fn extend_params(&mut self, p: impl IntoIterator<Item = (&'static str, String)>) {
1457        match self {
1458            PendingTxnResponse::Committed { params }
1459            | PendingTxnResponse::Rolledback { params } => params.extend(p),
1460        }
1461    }
1462}
1463
1464impl From<PendingTxnResponse> for ExecuteResponse {
1465    fn from(value: PendingTxnResponse) -> Self {
1466        match value {
1467            PendingTxnResponse::Committed { params } => {
1468                ExecuteResponse::TransactionCommitted { params }
1469            }
1470            PendingTxnResponse::Rolledback { params } => {
1471                ExecuteResponse::TransactionRolledBack { params }
1472            }
1473        }
1474    }
1475}
1476
1477#[derive(Debug)]
1478/// A pending read transaction waiting to be linearized along with metadata about it's state
1479pub struct PendingReadTxn {
1480    /// The transaction type
1481    txn: PendingRead,
1482    /// The timestamp context of the transaction.
1483    timestamp_context: TimestampContext,
1484    /// When we created this pending txn, when the transaction ends. Only used for metrics.
1485    created: Instant,
1486    /// Number of times we requeued the processing of this pending read txn.
1487    /// Requeueing is necessary if the time we executed the query is after the current oracle time;
1488    /// see [`Coordinator::message_linearize_reads`] for more details.
1489    num_requeues: u64,
1490    /// Telemetry context.
1491    otel_ctx: OpenTelemetryContext,
1492}
1493
1494impl PendingReadTxn {
1495    /// Return the timestamp context of the pending read transaction.
1496    pub fn timestamp_context(&self) -> &TimestampContext {
1497        &self.timestamp_context
1498    }
1499
1500    pub(crate) fn take_context(self) -> ExecuteContext {
1501        self.txn.take_context()
1502    }
1503}
1504
1505#[derive(Debug)]
1506/// A pending read transaction waiting to be linearized.
1507enum PendingRead {
1508    Read {
1509        /// The inner transaction.
1510        txn: PendingTxn,
1511    },
1512    ReadThenWrite {
1513        /// Context used to send a response back to the client.
1514        ctx: ExecuteContext,
1515        /// Channel used to alert the transaction that the read has been linearized and send back
1516        /// `ctx`.
1517        tx: oneshot::Sender<Option<ExecuteContext>>,
1518    },
1519}
1520
1521impl PendingRead {
1522    /// Alert the client that the read has been linearized.
1523    ///
1524    /// If it is necessary to finalize an execute, return the state necessary to do so
1525    /// (execution context and result)
1526    #[instrument(level = "debug")]
1527    pub fn finish(self) -> Option<(ExecuteContext, Result<ExecuteResponse, AdapterError>)> {
1528        match self {
1529            PendingRead::Read {
1530                txn:
1531                    PendingTxn {
1532                        mut ctx,
1533                        response,
1534                        action,
1535                    },
1536                ..
1537            } => {
1538                let changed = ctx.session_mut().vars_mut().end_transaction(action);
1539                // Append any parameters that changed to the response.
1540                let response = response.map(|mut r| {
1541                    r.extend_params(changed);
1542                    ExecuteResponse::from(r)
1543                });
1544
1545                Some((ctx, response))
1546            }
1547            PendingRead::ReadThenWrite { ctx, tx, .. } => {
1548                // Ignore errors if the caller has hung up.
1549                let _ = tx.send(Some(ctx));
1550                None
1551            }
1552        }
1553    }
1554
1555    fn label(&self) -> &'static str {
1556        match self {
1557            PendingRead::Read { .. } => "read",
1558            PendingRead::ReadThenWrite { .. } => "read_then_write",
1559        }
1560    }
1561
1562    pub(crate) fn take_context(self) -> ExecuteContext {
1563        match self {
1564            PendingRead::Read { txn, .. } => txn.ctx,
1565            PendingRead::ReadThenWrite { ctx, tx, .. } => {
1566                // Inform the transaction that we've taken their context.
1567                // Ignore errors if the caller has hung up.
1568                let _ = tx.send(None);
1569                ctx
1570            }
1571        }
1572    }
1573}
1574
1575/// State that the coordinator must process as part of retiring
1576/// command execution.  `ExecuteContextExtra::Default` is guaranteed
1577/// to produce a value that will cause the coordinator to do nothing, and
1578/// is intended for use by code that invokes the execution processing flow
1579/// (i.e., `sequence_plan`) without actually being a statement execution.
1580///
1581/// This is a pure data struct containing only the statement logging ID.
1582/// For auto-retire-on-drop behavior, use `ExecuteContextGuard` which wraps
1583/// this struct and owns the channel for sending retirement messages.
1584#[derive(Debug, Default)]
1585#[must_use]
1586pub struct ExecuteContextExtra {
1587    statement_uuid: Option<StatementLoggingId>,
1588}
1589
1590impl ExecuteContextExtra {
1591    pub(crate) fn new(statement_uuid: Option<StatementLoggingId>) -> Self {
1592        Self { statement_uuid }
1593    }
1594    pub fn is_trivial(&self) -> bool {
1595        self.statement_uuid.is_none()
1596    }
1597    pub fn contents(&self) -> Option<StatementLoggingId> {
1598        self.statement_uuid
1599    }
1600    /// Consume this extra and return the statement UUID for retirement.
1601    /// This should only be called from code that knows what to do to finish
1602    /// up logging based on the inner value.
1603    #[must_use]
1604    pub(crate) fn retire(self) -> Option<StatementLoggingId> {
1605        self.statement_uuid
1606    }
1607}
1608
1609/// A guard that wraps `ExecuteContextExtra` and owns a channel for sending
1610/// retirement messages to the coordinator.
1611///
1612/// If this guard is dropped with a `Some` `statement_uuid` in its inner
1613/// `ExecuteContextExtra`, the `Drop` implementation will automatically send a
1614/// `Message::RetireExecute` to log the statement ending.
1615/// This handles cases like connection drops where the context cannot be
1616/// explicitly retired.
1617/// See <https://github.com/MaterializeInc/database-issues/issues/7304>
1618#[derive(Debug)]
1619#[must_use]
1620pub struct ExecuteContextGuard {
1621    extra: ExecuteContextExtra,
1622    /// Channel for sending messages to the coordinator. Used for auto-retiring on drop.
1623    /// For `Default` instances, this is a dummy sender (receiver already dropped), so
1624    /// sends will fail silently - which is the desired behavior since Default instances
1625    /// should only be used for non-logged statements.
1626    coordinator_tx: mpsc::UnboundedSender<Message>,
1627}
1628
1629impl Default for ExecuteContextGuard {
1630    fn default() -> Self {
1631        // Create a dummy sender by immediately dropping the receiver.
1632        // Any send on this channel will fail silently, which is the desired
1633        // behavior for Default instances (non-logged statements).
1634        let (tx, _rx) = mpsc::unbounded_channel();
1635        Self {
1636            extra: ExecuteContextExtra::default(),
1637            coordinator_tx: tx,
1638        }
1639    }
1640}
1641
1642impl ExecuteContextGuard {
1643    pub(crate) fn new(
1644        statement_uuid: Option<StatementLoggingId>,
1645        coordinator_tx: mpsc::UnboundedSender<Message>,
1646    ) -> Self {
1647        Self {
1648            extra: ExecuteContextExtra::new(statement_uuid),
1649            coordinator_tx,
1650        }
1651    }
1652    pub fn is_trivial(&self) -> bool {
1653        self.extra.is_trivial()
1654    }
1655    pub fn contents(&self) -> Option<StatementLoggingId> {
1656        self.extra.contents()
1657    }
1658    /// Take responsibility for the contents.  This should only be
1659    /// called from code that knows what to do to finish up logging
1660    /// based on the inner value.
1661    ///
1662    /// Returns the inner `ExecuteContextExtra`, consuming the guard without
1663    /// triggering the auto-retire behavior.
1664    pub(crate) fn defuse(mut self) -> ExecuteContextExtra {
1665        // Taking statement_uuid prevents the Drop impl from sending a retire message
1666        std::mem::take(&mut self.extra)
1667    }
1668}
1669
1670impl Drop for ExecuteContextGuard {
1671    fn drop(&mut self) {
1672        if let Some(statement_uuid) = self.extra.statement_uuid.take() {
1673            // Auto-retire since the guard was dropped without explicit retirement (likely due
1674            // to connection drop).
1675            let msg = Message::RetireExecute {
1676                data: ExecuteContextExtra {
1677                    statement_uuid: Some(statement_uuid),
1678                },
1679                otel_ctx: OpenTelemetryContext::obtain(),
1680                reason: StatementEndedExecutionReason::Aborted,
1681            };
1682            // Send may fail for Default instances (dummy sender), which is fine since
1683            // Default instances should only be used for non-logged statements.
1684            let _ = self.coordinator_tx.send(msg);
1685        }
1686    }
1687}
1688
1689/// Carries the session and statement state needed to retire an execution.
1690///
1691/// Dropping an unretired context fails the client synchronously. Shutdown can drop contexts from
1692/// task queues, where spawning response-barrier work is no longer safe.
1693#[derive(Debug)]
1694pub struct ExecuteContext {
1695    // `None` only after `retire`/`into_parts` consumed the context.
1696    inner: Option<Box<ExecuteContextInner>>,
1697}
1698
1699impl std::ops::Deref for ExecuteContext {
1700    type Target = ExecuteContextInner;
1701    fn deref(&self) -> &Self::Target {
1702        self.inner.as_ref().expect("only consumed by value")
1703    }
1704}
1705
1706impl std::ops::DerefMut for ExecuteContext {
1707    fn deref_mut(&mut self) -> &mut Self::Target {
1708        self.inner.as_mut().expect("only consumed by value")
1709    }
1710}
1711
1712impl Drop for ExecuteContext {
1713    fn drop(&mut self) {
1714        let Some(inner) = self.inner.take() else {
1715            return;
1716        };
1717        // Destructors cannot spawn response-barrier tasks during runtime shutdown. Send the error
1718        // synchronously and let the statement guard report retirement.
1719        tracing::warn!("execute context dropped without retirement, failing the client");
1720        let ExecuteContextInner { tx, session, .. } = *inner;
1721        tx.send(
1722            Err(AdapterError::Internal(
1723                "statement execution abandoned, outcome unknown (server shutting down)".into(),
1724            )),
1725            session,
1726        );
1727    }
1728}
1729
1730#[derive(Derivative)]
1731#[derivative(Debug)]
1732pub struct ExecuteContextInner {
1733    tx: ClientTransmitter<ExecuteResponse>,
1734    internal_cmd_tx: mpsc::UnboundedSender<Message>,
1735    session: Session,
1736    extra: ExecuteContextGuard,
1737    #[derivative(Debug = "ignore")]
1738    response_barriers: Vec<BuiltinTableAppendNotify>,
1739}
1740
1741impl ExecuteContext {
1742    pub fn session(&self) -> &Session {
1743        &self.session
1744    }
1745
1746    pub fn session_mut(&mut self) -> &mut Session {
1747        &mut self.session
1748    }
1749
1750    pub fn tx(&self) -> &ClientTransmitter<ExecuteResponse> {
1751        &self.tx
1752    }
1753
1754    pub fn tx_mut(&mut self) -> &mut ClientTransmitter<ExecuteResponse> {
1755        &mut self.tx
1756    }
1757
1758    pub fn from_parts(
1759        tx: ClientTransmitter<ExecuteResponse>,
1760        internal_cmd_tx: mpsc::UnboundedSender<Message>,
1761        session: Session,
1762        extra: ExecuteContextGuard,
1763    ) -> Self {
1764        Self::from_parts_with_response_barriers(tx, internal_cmd_tx, session, extra, Vec::new())
1765    }
1766
1767    pub fn from_parts_with_response_barriers(
1768        tx: ClientTransmitter<ExecuteResponse>,
1769        internal_cmd_tx: mpsc::UnboundedSender<Message>,
1770        session: Session,
1771        extra: ExecuteContextGuard,
1772        response_barriers: Vec<BuiltinTableAppendNotify>,
1773    ) -> Self {
1774        Self {
1775            inner: Some(
1776                ExecuteContextInner {
1777                    tx,
1778                    session,
1779                    extra,
1780                    response_barriers,
1781                    internal_cmd_tx,
1782                }
1783                .into(),
1784            ),
1785        }
1786    }
1787
1788    /// By calling this function, the caller takes responsibility for
1789    /// dealing with the instance of `ExecuteContextGuard`. This is
1790    /// intended to support protocols (like `COPY FROM`) that involve
1791    /// multiple passes of sending the session back and forth between
1792    /// the coordinator and the pgwire layer. As part of any such
1793    /// protocol, we must ensure that the `ExecuteContextGuard`
1794    /// (possibly wrapped in a new `ExecuteContext`) is passed back to the coordinator for
1795    /// eventual retirement. The returned response barriers must stay attached
1796    /// to the user-visible response path.
1797    ///
1798    /// The returned parts lose the `Drop` backstop that answers the client on shutdown, so they
1799    /// must not be held across an await point. A bare `ClientTransmitter` panics when dropped
1800    /// unsent.
1801    pub fn into_parts(
1802        mut self,
1803    ) -> (
1804        ClientTransmitter<ExecuteResponse>,
1805        mpsc::UnboundedSender<Message>,
1806        Session,
1807        ExecuteContextGuard,
1808        Vec<BuiltinTableAppendNotify>,
1809    ) {
1810        let ExecuteContextInner {
1811            tx,
1812            internal_cmd_tx,
1813            session,
1814            extra,
1815            response_barriers,
1816        } = *self.inner.take().expect("only consumed by value");
1817        (tx, internal_cmd_tx, session, extra, response_barriers)
1818    }
1819
1820    /// Retire the execution, by sending a message to the coordinator.
1821    #[instrument(level = "debug")]
1822    pub fn retire(mut self, result: Result<ExecuteResponse, AdapterError>) {
1823        let response_barriers = std::mem::take(&mut self.response_barriers);
1824        if response_barriers.is_empty() {
1825            let (tx, internal_cmd_tx, session, extra, _) = self.into_parts();
1826            retire_execution_context(tx, internal_cmd_tx, session, extra, result);
1827            return;
1828        }
1829        // Keep `self` intact across the wait: if shutdown drops this task, the context's `Drop`
1830        // backstop answers the client. Barriers are empty on re-entry, so this terminates.
1831        spawn(
1832            || "execute_context::retire_after_response_barriers",
1833            async move {
1834                for barrier in response_barriers {
1835                    barrier.await;
1836                }
1837                self.retire(result);
1838            },
1839        );
1840    }
1841
1842    /// Delays sending this statement's response until `barrier` resolves.
1843    pub(crate) fn delay_response_until(&mut self, barrier: BuiltinTableAppendCompletion) {
1844        self.response_barriers.push(barrier.into_notify());
1845    }
1846
1847    pub fn extra(&self) -> &ExecuteContextGuard {
1848        &self.extra
1849    }
1850
1851    pub fn extra_mut(&mut self) -> &mut ExecuteContextGuard {
1852        &mut self.extra
1853    }
1854}
1855
1856fn retire_execution_context(
1857    tx: ClientTransmitter<ExecuteResponse>,
1858    internal_cmd_tx: mpsc::UnboundedSender<Message>,
1859    session: Session,
1860    extra: ExecuteContextGuard,
1861    result: Result<ExecuteResponse, AdapterError>,
1862) {
1863    let reason = if extra.is_trivial() {
1864        None
1865    } else {
1866        Some((&result).into())
1867    };
1868    tx.send(result, session);
1869    if let Some(reason) = reason {
1870        let extra = extra.defuse();
1871        if let Err(e) = internal_cmd_tx.send(Message::RetireExecute {
1872            otel_ctx: OpenTelemetryContext::obtain(),
1873            data: extra,
1874            reason,
1875        }) {
1876            warn!("internal_cmd_rx dropped before we could send: {:?}", e);
1877        }
1878    }
1879}
1880
1881#[derive(Debug)]
1882struct ClusterReplicaStatuses(
1883    BTreeMap<ClusterId, BTreeMap<ReplicaId, BTreeMap<ProcessId, ClusterReplicaProcessStatus>>>,
1884);
1885
1886impl ClusterReplicaStatuses {
1887    pub(crate) fn new() -> ClusterReplicaStatuses {
1888        ClusterReplicaStatuses(BTreeMap::new())
1889    }
1890
1891    /// Initializes the statuses of the specified cluster.
1892    ///
1893    /// Panics if the cluster statuses are already initialized.
1894    pub(crate) fn initialize_cluster_statuses(&mut self, cluster_id: ClusterId) {
1895        let prev = self.0.insert(cluster_id, BTreeMap::new());
1896        assert_eq!(
1897            prev, None,
1898            "cluster {cluster_id} statuses already initialized"
1899        );
1900    }
1901
1902    /// Initializes the statuses of the specified cluster replica.
1903    ///
1904    /// Panics if the cluster replica statuses are already initialized.
1905    pub(crate) fn initialize_cluster_replica_statuses(
1906        &mut self,
1907        cluster_id: ClusterId,
1908        replica_id: ReplicaId,
1909        num_processes: usize,
1910        time: DateTime<Utc>,
1911    ) {
1912        tracing::info!(
1913            ?cluster_id,
1914            ?replica_id,
1915            ?time,
1916            "initializing cluster replica status"
1917        );
1918        let replica_statuses = self.0.entry(cluster_id).or_default();
1919        let process_statuses = (0..num_processes)
1920            .map(|process_id| {
1921                let status = ClusterReplicaProcessStatus {
1922                    status: ClusterStatus::Offline(Some(OfflineReason::Initializing)),
1923                    restart_count: 0,
1924                    time: time.clone(),
1925                };
1926                (u64::cast_from(process_id), status)
1927            })
1928            .collect();
1929        let prev = replica_statuses.insert(replica_id, process_statuses);
1930        assert_none!(
1931            prev,
1932            "cluster replica {cluster_id}.{replica_id} statuses already initialized"
1933        );
1934    }
1935
1936    /// Removes the statuses of the specified cluster.
1937    ///
1938    /// Panics if the cluster does not exist.
1939    pub(crate) fn remove_cluster_statuses(
1940        &mut self,
1941        cluster_id: &ClusterId,
1942    ) -> BTreeMap<ReplicaId, BTreeMap<ProcessId, ClusterReplicaProcessStatus>> {
1943        let prev = self.0.remove(cluster_id);
1944        prev.unwrap_or_else(|| panic!("unknown cluster: {cluster_id}"))
1945    }
1946
1947    /// Removes the statuses of the specified cluster replica.
1948    ///
1949    /// Panics if the cluster or replica does not exist.
1950    pub(crate) fn remove_cluster_replica_statuses(
1951        &mut self,
1952        cluster_id: &ClusterId,
1953        replica_id: &ReplicaId,
1954    ) -> BTreeMap<ProcessId, ClusterReplicaProcessStatus> {
1955        let replica_statuses = self
1956            .0
1957            .get_mut(cluster_id)
1958            .unwrap_or_else(|| panic!("unknown cluster: {cluster_id}"));
1959        let prev = replica_statuses.remove(replica_id);
1960        prev.unwrap_or_else(|| panic!("unknown cluster replica: {cluster_id}.{replica_id}"))
1961    }
1962
1963    /// Inserts or updates the status of the specified cluster replica process.
1964    ///
1965    /// Panics if the cluster or replica does not exist.
1966    pub(crate) fn ensure_cluster_status(
1967        &mut self,
1968        cluster_id: ClusterId,
1969        replica_id: ReplicaId,
1970        process_id: ProcessId,
1971        status: ClusterReplicaProcessStatus,
1972    ) {
1973        let replica_statuses = self
1974            .0
1975            .get_mut(&cluster_id)
1976            .unwrap_or_else(|| panic!("unknown cluster: {cluster_id}"))
1977            .get_mut(&replica_id)
1978            .unwrap_or_else(|| panic!("unknown cluster replica: {cluster_id}.{replica_id}"));
1979        replica_statuses.insert(process_id, status);
1980    }
1981
1982    /// Computes the status of the cluster replica as a whole.
1983    ///
1984    /// Panics if `cluster_id` or `replica_id` don't exist.
1985    pub fn get_cluster_replica_status(
1986        &self,
1987        cluster_id: ClusterId,
1988        replica_id: ReplicaId,
1989    ) -> ClusterStatus {
1990        let process_status = self.get_cluster_replica_statuses(cluster_id, replica_id);
1991        Self::cluster_replica_status(process_status)
1992    }
1993
1994    /// Computes the status of the cluster replica as a whole.
1995    pub fn cluster_replica_status(
1996        process_status: &BTreeMap<ProcessId, ClusterReplicaProcessStatus>,
1997    ) -> ClusterStatus {
1998        process_status
1999            .values()
2000            .fold(ClusterStatus::Online, |s, p| match (s, p.status) {
2001                (ClusterStatus::Online, ClusterStatus::Online) => ClusterStatus::Online,
2002                (x, y) => {
2003                    let reason_x = match x {
2004                        ClusterStatus::Offline(reason) => reason,
2005                        ClusterStatus::Online => None,
2006                    };
2007                    let reason_y = match y {
2008                        ClusterStatus::Offline(reason) => reason,
2009                        ClusterStatus::Online => None,
2010                    };
2011                    // Arbitrarily pick the first known not-ready reason.
2012                    ClusterStatus::Offline(reason_x.or(reason_y))
2013                }
2014            })
2015    }
2016
2017    /// Gets the statuses of the given cluster replica.
2018    ///
2019    /// Panics if the cluster or replica does not exist
2020    pub(crate) fn get_cluster_replica_statuses(
2021        &self,
2022        cluster_id: ClusterId,
2023        replica_id: ReplicaId,
2024    ) -> &BTreeMap<ProcessId, ClusterReplicaProcessStatus> {
2025        self.try_get_cluster_replica_statuses(cluster_id, replica_id)
2026            .unwrap_or_else(|| panic!("unknown cluster replica: {cluster_id}.{replica_id}"))
2027    }
2028
2029    /// Gets the statuses of the given cluster replica.
2030    pub(crate) fn try_get_cluster_replica_statuses(
2031        &self,
2032        cluster_id: ClusterId,
2033        replica_id: ReplicaId,
2034    ) -> Option<&BTreeMap<ProcessId, ClusterReplicaProcessStatus>> {
2035        self.try_get_cluster_statuses(cluster_id)
2036            .and_then(|statuses| statuses.get(&replica_id))
2037    }
2038
2039    /// Gets the statuses of the given cluster.
2040    pub(crate) fn try_get_cluster_statuses(
2041        &self,
2042        cluster_id: ClusterId,
2043    ) -> Option<&BTreeMap<ReplicaId, BTreeMap<ProcessId, ClusterReplicaProcessStatus>>> {
2044        self.0.get(&cluster_id)
2045    }
2046}
2047
2048/// Glues the external world to the Timely workers.
2049#[derive(Derivative)]
2050#[derivative(Debug)]
2051pub struct Coordinator {
2052    /// The controller for the storage and compute layers.
2053    #[derivative(Debug = "ignore")]
2054    controller: mz_controller::Controller,
2055    /// The catalog in an Arc suitable for readonly references. The Arc allows
2056    /// us to hand out cheap copies of the catalog to functions that can use it
2057    /// off of the main coordinator thread. If the coordinator needs to mutate
2058    /// the catalog, call [`Self::catalog_mut`], which will clone this struct member,
2059    /// allowing it to be mutated here while the other off-thread references can
2060    /// read their catalog as long as needed. In the future we would like this
2061    /// to be a pTVC, but for now this is sufficient.
2062    catalog: Arc<Catalog>,
2063
2064    /// A client for persist. Initially, this is only used for reading stashed
2065    /// peek responses out of batches.
2066    persist_client: PersistClient,
2067
2068    /// Channel to manage internal commands from the coordinator to itself.
2069    internal_cmd_tx: mpsc::UnboundedSender<Message>,
2070    /// Notification that triggers a group commit.
2071    group_commit_tx: appends::GroupCommitNotifier,
2072    /// Wakes the cluster controller task to reconcile immediately instead of
2073    /// waiting out its tick interval. Notified after catalog transactions that
2074    /// change durable cluster state.
2075    reconcile_now: Arc<Notify>,
2076    group_committer_tx: mpsc::UnboundedSender<appends::TableWriteCmd>,
2077
2078    /// Channel for strict serializable reads ready to commit.
2079    strict_serializable_reads_tx: mpsc::UnboundedSender<(ConnectionId, PendingReadTxn)>,
2080
2081    /// Signals that pending strict serializable reads should be re-checked
2082    /// because the timestamp oracle may have advanced. Awaited below group commit
2083    /// in [`Coordinator::serve`]; see that branch for the ordering rationale.
2084    linearize_reads_notify: Arc<Notify>,
2085
2086    /// Mechanism for totally ordering write and read timestamps, so that all reads
2087    /// reflect exactly the set of writes that precede them, and no writes that follow.
2088    global_timelines: BTreeMap<Timeline, TimelineState>,
2089
2090    /// A generator for transient [`GlobalId`]s, shareable with other threads.
2091    transient_id_gen: Arc<TransientIdGen>,
2092    /// A map from connection ID to metadata about that connection for all
2093    /// active connections.
2094    active_conns: BTreeMap<ConnectionId, ConnMeta>,
2095
2096    /// For each transaction, the read holds taken to support any performed reads.
2097    ///
2098    /// Upon completing a transaction, these read holds should be dropped.
2099    txn_read_holds: BTreeMap<ConnectionId, read_policy::ReadHolds>,
2100
2101    /// Access to the peek fields should be restricted to methods in the [`peek`] API.
2102    /// A map from pending peek ids to the queue into which responses are sent, and
2103    /// the connection id of the client that initiated the peek.
2104    pending_peeks: BTreeMap<Uuid, PendingPeek>,
2105    /// A map from client connection ids to a set of all pending peeks for that client.
2106    client_pending_peeks: BTreeMap<ConnectionId, BTreeMap<Uuid, ClusterId>>,
2107
2108    /// A map from client connection ids to pending linearize read transaction.
2109    pending_linearize_read_txns: BTreeMap<ConnectionId, PendingReadTxn>,
2110
2111    /// A map from the compute sink ID to it's state description.
2112    active_compute_sinks: BTreeMap<GlobalId, ActiveComputeSink>,
2113    /// A map from active webhooks to their invalidation handle.
2114    active_webhooks: BTreeMap<CatalogItemId, WebhookAppenderInvalidator>,
2115    /// A map of active `COPY FROM` statements. The Coordinator waits for `clusterd`
2116    /// to stage Batches in Persist that we will then link into the shard.
2117    active_copies: BTreeMap<ConnectionId, ActiveCopyFrom>,
2118
2119    /// Connection-scoped cancellation watches.
2120    ///
2121    /// Each entry is a watch channel whose value is `false` until cancellation
2122    /// is requested for that connection, at which point it is set to `true`.
2123    ///
2124    /// Consumers install these watches while they have cancellable work in
2125    /// flight, always as a fresh channel, so nobody can observe a cancellation
2126    /// aimed at an earlier statement. An entry is removed when a statement
2127    /// starts, when a stage runs uncancelable, and when the connection's state
2128    /// is cleared.
2129    connection_cancel_watches: BTreeMap<ConnectionId, (watch::Sender<bool>, watch::Receiver<bool>)>,
2130    /// Active introspection subscribes.
2131    introspection_subscribes: BTreeMap<GlobalId, IntrospectionSubscribe>,
2132    /// The last replica visited by the sequential hydration-history sweep.
2133    hydration_history_replica_cursor: Option<ReplicaId>,
2134    /// Hydration-history sweep owned by the coordinator while one is in flight.
2135    hydration_history_sweep: Option<AbortOnDropHandle<()>>,
2136    /// The curated metric sinks installed on each replica.
2137    ///
2138    /// Keyed replica-first so a replica's installs form one contiguous range: teardown on replica
2139    /// drop is the only lookup that is not by exact key.
2140    metric_sinks: BTreeMap<(ReplicaId, &'static str), InstalledMetricSink>,
2141    /// Curated metric-sink plans, cached per definition so each is planned once rather than once
2142    /// per replica. See [`Coordinator::plan_metric_sink`].
2143    metric_sink_plans: BTreeMap<&'static str, PlannedMetricSink>,
2144
2145    /// Locks that grant access to a specific object, populated lazily as objects are written to.
2146    write_locks: BTreeMap<CatalogItemId, Arc<tokio::sync::Mutex<()>>>,
2147    /// Plans that are currently deferred and waiting on a write lock.
2148    deferred_write_ops: BTreeMap<ConnectionId, DeferredOp>,
2149
2150    /// Pending writes waiting for a group commit.
2151    pending_writes: Vec<PendingWriteTxn>,
2152
2153    /// Semaphore to limit concurrent OCC (optimistic concurrency control)
2154    /// read-then-write operations.
2155    ///
2156    /// Each operation maintains a subscribe that continually receives and
2157    /// consolidates updates. With N concurrent loops, every successful write
2158    /// forces the other N-1 to redo work, so total work scales as `O(n^2)`.
2159    /// The semaphore caps concurrency to keep that bounded.
2160    ///
2161    /// NOTE: The number of permits is read from `max_concurrent_occ_writes` at
2162    /// coordinator startup. Runtime changes require an `environmentd` restart.
2163    occ_write_semaphore: Arc<Semaphore>,
2164
2165    /// Whether frontend OCC read-then-write is enabled. Read once at startup
2166    /// from the `FRONTEND_READ_THEN_WRITE` dyncfg and fixed for the lifetime of
2167    /// this process. See the module-level docs on `frontend_read_then_write`
2168    /// for why mixed-mode operation is not allowed.
2169    frontend_read_then_write_enabled: bool,
2170
2171    /// For the realtime timeline, an explicit SELECT or INSERT on a table will bump the
2172    /// table's timestamps, but there are cases where timestamps are not bumped but
2173    /// we expect the closed timestamps to advance (`AS OF X`, SUBSCRIBing views over
2174    /// RT sources and tables). To address these, spawn a task that forces table
2175    /// timestamps to close on a regular interval. This roughly tracks the behavior
2176    /// of realtime sources that close off timestamps on an interval.
2177    ///
2178    /// For non-realtime timelines, nothing pushes the timestamps forward, so we must do
2179    /// it manually.
2180    advance_timelines_interval: Interval,
2181
2182    /// Serialized DDL. DDL must be serialized because:
2183    /// - Many of them do off-thread work and need to verify the catalog is in a valid state, but
2184    ///   [`PlanValidity`] does not currently support tracking all changes. Doing that correctly
2185    ///   seems to be more difficult than it's worth, so we would instead re-plan and re-sequence
2186    ///   the statements.
2187    /// - Re-planning a statement is hard because Coordinator and Session state is mutated at
2188    ///   various points, and we would need to correctly reset those changes before re-planning and
2189    ///   re-sequencing.
2190    serialized_ddl: LockedVecDeque<DeferredPlanStatement>,
2191
2192    /// Handle to secret manager that can create and delete secrets from
2193    /// an arbitrary secret storage engine.
2194    secrets_controller: Arc<dyn SecretsController>,
2195    /// A secrets reader than maintains an in-memory cache, where values have a set TTL.
2196    caching_secrets_reader: CachingSecretsReader,
2197
2198    /// Handle to a manager that can create and delete kubernetes resources
2199    /// (ie: VpcEndpoint objects)
2200    cloud_resource_controller: Option<Arc<dyn CloudResourceController>>,
2201
2202    /// Persist client for fetching storage metadata such as size metrics.
2203    storage_usage_client: StorageUsageClient,
2204    /// The interval at which to collect storage usage information.
2205    storage_usage_collection_interval: Duration,
2206
2207    /// Segment analytics client.
2208    #[derivative(Debug = "ignore")]
2209    segment_client: Option<mz_segment::Client>,
2210
2211    /// Coordinator metrics.
2212    metrics: Metrics,
2213    /// Optimizer metrics.
2214    optimizer_metrics: OptimizerMetrics,
2215
2216    /// Tracing handle.
2217    tracing_handle: TracingHandle,
2218
2219    /// Data used by the statement logging feature.
2220    statement_logging: StatementLogging,
2221
2222    /// Limit for how many concurrent webhook requests we allow.
2223    webhook_concurrency_limit: WebhookConcurrencyLimiter,
2224
2225    /// Optional config for the timestamp oracle. This is _required_ when
2226    /// a timestamp oracle backend is configured.
2227    timestamp_oracle_config: Option<TimestampOracleConfig>,
2228
2229    /// When doing 0dt upgrades/in read-only mode, periodically ask all known
2230    /// clusters/collections whether they are caught up.
2231    caught_up_check_interval: Interval,
2232
2233    /// Context needed to check whether all clusters/collections have caught up.
2234    /// Only used during 0dt deployment, while in read-only mode.
2235    caught_up_check: Option<CaughtUpCheckContext>,
2236
2237    /// The metrics registry, handed to the catalog info-metrics background task
2238    /// so it can register and own its `*_info` series.
2239    catalog_info_metrics_registry: MetricsRegistry,
2240
2241    /// The shared system-parameter frontend, installed by the sync loop once it
2242    /// initializes (and re-installed on reconnect). `None` until then, for
2243    /// example before LaunchDarkly connects, where a newly-created object
2244    /// resolves to the environment-wide value (the cold-cache fallback). Used to
2245    /// resolve a new cluster's or replica's scoped overrides synchronously at
2246    /// create time, so its first plan or first controller configuration is
2247    /// correct rather than waiting for the next sync tick. See the scoped
2248    /// feature flags design.
2249    scoped_frontend: Option<Arc<SystemParameterFrontend>>,
2250
2251    /// Tracks the state associated with the currently installed watchsets.
2252    installed_watch_sets: BTreeMap<WatchSetId, (ConnectionId, WatchSetResponse)>,
2253
2254    /// Tracks the currently installed watchsets for each connection.
2255    connection_watch_sets: BTreeMap<ConnectionId, BTreeSet<WatchSetId>>,
2256
2257    /// Tracks the statuses of all cluster replicas.
2258    cluster_replica_statuses: ClusterReplicaStatuses,
2259
2260    /// Whether or not to start controllers in read-only mode. This is only
2261    /// meant for use during development of read-only clusters and 0dt upgrades
2262    /// and should go away once we have proper orchestration during upgrades.
2263    read_only_controllers: bool,
2264
2265    /// Updates to builtin tables that are being buffered while we are in
2266    /// read-only mode. We apply these all at once when coming out of read-only
2267    /// mode.
2268    ///
2269    /// This is a `Some` while in read-only mode and will be replaced by a
2270    /// `None` when we transition out of read-only mode and write out any
2271    /// buffered updates.
2272    buffered_builtin_table_updates: Option<Vec<BuiltinTableUpdate>>,
2273
2274    license_key: ValidatedLicenseKey,
2275
2276    /// Pre-allocated pool of user IDs to amortize persist writes across DDL operations.
2277    user_id_pool: IdPool,
2278}
2279
2280impl Coordinator {
2281    /// Persists the scoped system-parameter working copy and reconciles it into
2282    /// the per-scope resolution boundaries.
2283    ///
2284    /// The system-parameter sync loop and the create-time fold
2285    /// (`scoped_overrides_create_op`, folded into the create transaction) are the
2286    /// only writers, both serialized on the coordinator loop. The diff is
2287    /// persisted to the
2288    /// durable cache (so values survive an `environmentd` restart and an LD
2289    /// outage) via `Op::UpdateScopedSystemParameters`, which also updates the
2290    /// in-memory working copy in [`CatalogState`] and the
2291    /// `mz_cluster_system_parameters` / `mz_replica_system_parameters`
2292    /// introspection relations. The `replica`-scoped overrides reach the compute
2293    /// controller's per-replica dyncfg layer through the catalog implication for
2294    /// the persisted change. The `cluster`-scoped layer is resolved at plan time
2295    /// via [`CatalogState::cluster_scoped_optimizer_overrides`].
2296    ///
2297    /// [`CatalogState`]: crate::catalog::CatalogState
2298    /// [`CatalogState::cluster_scoped_optimizer_overrides`]: crate::catalog::CatalogState::cluster_scoped_optimizer_overrides
2299    pub(crate) async fn reconcile_scoped_system_parameters(
2300        &mut self,
2301        scoped: ScopedParameters,
2302        prune_scope: ScopedParametersScope,
2303    ) {
2304        // Nothing changed: skip the durable write. This is the common case on
2305        // most sync ticks.
2306        if self.catalog().state().scoped_system_parameters() == &scoped {
2307            return;
2308        }
2309
2310        // Persist the diff and update the in-memory working copy + introspection
2311        // through the catalog transaction, serialized on the coordinator loop
2312        // with the create-time fold. The replica-scoped
2313        // controller push is derived from this transaction's diff by the catalog
2314        // implication. `prune_scope` bounds removals to the evaluated objects, so
2315        // a concurrently-created object's override is not wiped. Best-effort: a
2316        // failure here is logged and retried on the next sync tick.
2317        if let Err(e) = self
2318            .catalog_transact(
2319                None,
2320                vec![crate::catalog::Op::UpdateScopedSystemParameters {
2321                    scoped,
2322                    prune_scope,
2323                }],
2324            )
2325            .await
2326        {
2327            tracing::warn!("failed to persist scoped system parameters: {e}");
2328        }
2329    }
2330
2331    /// Evaluates scoped overrides for objects created by `ops` and returns an
2332    /// [`Op::UpdateScopedSystemParameters`] to fold into the same transaction.
2333    ///
2334    /// The objects are not yet in the catalog, so this derives their contexts
2335    /// from concrete create ops and pre-allocated ids. Centralizing the fold
2336    /// here makes create-time configuration an invariant of coordinator-applied
2337    /// catalog ops, independent of which component produced them. The committed
2338    /// diff drives the replica-scoped controller push before `create_replica`.
2339    /// Render-frozen flags make a later push too late.
2340    ///
2341    /// Returns `None` when no scoped object is created or the shared frontend is
2342    /// not yet installed. An installed frontend produces an op even when no
2343    /// override applies, so a final DDL-transaction evaluation can clear a value
2344    /// staged by an earlier statement. The periodic sync loop remains the
2345    /// authoritative full-state reconciler.
2346    ///
2347    /// [`Op::UpdateScopedSystemParameters`]: crate::catalog::Op::UpdateScopedSystemParameters
2348    fn scoped_overrides_create_op(&self, ops: &[crate::catalog::Op]) -> Option<crate::catalog::Op> {
2349        let mut created_clusters = BTreeMap::new();
2350        let mut clusters = Vec::new();
2351        for op in ops {
2352            let crate::catalog::Op::CreateCluster { id, name, .. } = op else {
2353                continue;
2354            };
2355            let cluster = ClusterScopeContext {
2356                id: id.to_string(),
2357                name: name.clone(),
2358                is_builtin: id.is_system(),
2359            };
2360            created_clusters.insert(*id, cluster.clone());
2361            clusters.push(ClusterEvalContext {
2362                cluster_id: *id,
2363                cluster,
2364            });
2365        }
2366
2367        let mut replicas = Vec::new();
2368        for op in ops {
2369            let crate::catalog::Op::CreateClusterReplica {
2370                cluster_id,
2371                replica_id,
2372                name,
2373                config,
2374                ..
2375            } = op
2376            else {
2377                continue;
2378            };
2379            let ReplicaLocation::Managed(location) = &config.location else {
2380                continue;
2381            };
2382            let Some(cluster) = created_clusters.get(cluster_id).cloned().or_else(|| {
2383                self.catalog()
2384                    .try_get_cluster(*cluster_id)
2385                    .map(|cluster| ClusterScopeContext {
2386                        id: cluster_id.to_string(),
2387                        name: cluster.name.clone(),
2388                        is_builtin: cluster_id.is_system(),
2389                    })
2390            }) else {
2391                continue;
2392            };
2393            replicas.push(ReplicaEvalContext {
2394                cluster_id: *cluster_id,
2395                replica_id: *replica_id,
2396                replica: ReplicaScopeContext {
2397                    id: replica_id.to_string(),
2398                    name: name.clone(),
2399                    is_builtin: cluster_id.is_system(),
2400                    size: location.size.clone(),
2401                    size_family: location.allocation.family().to_string(),
2402                    cluster_id: cluster_id.to_string(),
2403                    cluster_name: cluster.name.clone(),
2404                },
2405                cluster,
2406            });
2407        }
2408
2409        if clusters.is_empty() && replicas.is_empty() {
2410            return None;
2411        }
2412        let frontend = self.scoped_frontend.clone()?;
2413        let catalog = self.catalog();
2414        let system_config = catalog.system_config();
2415
2416        // Partition the synced parameters by scope class, as the sync loop does,
2417        // so we evaluate exactly the flags in use at each scope.
2418        let replica_param_names: Vec<&'static str> = system_config
2419            .iter_synced()
2420            .filter(|var| var.scope() == ParameterScope::Replica)
2421            .map(|var| var.name())
2422            .collect();
2423        let cluster_param_names: Vec<&'static str> = system_config
2424            .iter_synced()
2425            .filter(|var| var.scope() == ParameterScope::Cluster)
2426            .map(|var| var.name())
2427            .collect();
2428
2429        let params = SynchronizedParameters::new(system_config.clone());
2430        let mut evaluated = ScopedParameters::default();
2431        if !cluster_param_names.is_empty() && !clusters.is_empty() {
2432            evaluated.cluster =
2433                frontend.pull_cluster_overrides(&params, &cluster_param_names, &clusters);
2434        }
2435        if !replica_param_names.is_empty() && !replicas.is_empty() {
2436            evaluated.replica =
2437                frontend.pull_replica_overrides(&params, &replica_param_names, &replicas);
2438        }
2439        // Prune only within the objects this transaction creates. A later
2440        // statement in a DDL transaction can replace an earlier folded value,
2441        // but this never touches an unrelated object's override.
2442        let prune_scope = ScopedParametersScope {
2443            clusters: clusters.iter().map(|cluster| cluster.cluster_id).collect(),
2444            replicas: replicas.iter().map(|replica| replica.replica_id).collect(),
2445        };
2446        Some(crate::catalog::Op::UpdateScopedSystemParameters {
2447            scoped: evaluated,
2448            prune_scope,
2449        })
2450    }
2451
2452    /// Renders the replica-local scoped overrides in the catalog working copy as
2453    /// per-replica [`ConfigUpdates`], grouped by cluster.
2454    ///
2455    /// Sparse: only replicas with an override are present. Parameters that are
2456    /// not dyncfgs are skipped, as are values that fail to parse.
2457    pub(crate) fn replica_dyncfg_overrides(
2458        &self,
2459    ) -> BTreeMap<ComputeInstanceId, BTreeMap<ReplicaId, ConfigUpdates>> {
2460        let replica_overrides = &self.catalog().state().scoped_system_parameters().replica;
2461
2462        let dyncfgs = self.catalog().system_config().dyncfgs();
2463        let mut instance_overrides: BTreeMap<
2464            ComputeInstanceId,
2465            BTreeMap<ReplicaId, ConfigUpdates>,
2466        > = BTreeMap::new();
2467        for cluster in self.catalog().clusters() {
2468            for replica in cluster.replicas() {
2469                let Some(values) = replica_overrides.get(&replica.replica_id) else {
2470                    continue;
2471                };
2472                let mut updates = ConfigUpdates::default();
2473                for (name, value) in values {
2474                    let Some(entry) = dyncfgs.entry(name) else {
2475                        // A replica-local parameter that is not a dyncfg has no
2476                        // per-replica realization, so skip it.
2477                        continue;
2478                    };
2479                    match entry.parse_val(value) {
2480                        Ok(val) => updates.add_dynamic(name, val),
2481                        Err(e) => {
2482                            tracing::warn!(%name, %value, "cannot parse scoped override: {e}")
2483                        }
2484                    }
2485                }
2486                if !updates.updates.is_empty() {
2487                    instance_overrides
2488                        .entry(cluster.id)
2489                        .or_default()
2490                        .insert(replica.replica_id, updates);
2491                }
2492            }
2493        }
2494
2495        instance_overrides
2496    }
2497
2498    /// Resolves the replica-local scoped overrides from the catalog working copy
2499    /// into the controllers' per-replica dyncfg layers, then re-pushes the
2500    /// environment-wide configuration so replicas observe the new values.
2501    /// Driven by the catalog implication for replica-scoped configuration
2502    /// changes, and called once on bootstrap.
2503    pub(crate) fn push_replica_dyncfg_overrides(&mut self) {
2504        let instance_overrides = self.replica_dyncfg_overrides();
2505
2506        // Both controllers carry a per-replica dyncfg layer, because the two
2507        // protocols realize configs in different worker `ConfigSet`s on
2508        // `clusterd`. The compute worker's `handle_update_configuration`
2509        // applies the pushed dyncfg updates to compute's own worker
2510        // `ConfigSet`, to the shared persist client `ConfigSet`
2511        // (`persist_clients.cfg()`) that the co-located storage server reads
2512        // from the same `Arc`, and to `mz_metrics`, which covers
2513        // persist-backed and process-global configs such as persist client
2514        // tuning and `lgalloc`. Configs realized from the storage worker's own
2515        // `ConfigSet` (read in its `UpdateConfiguration` handler) are reached
2516        // only by the storage controller's layer. A third class is not pushed
2517        // to a running replica at all but baked into its process configuration
2518        // when the controller provisions it, which is why the overrides also go
2519        // to the outer controller.
2520        self.controller
2521            .update_replica_dyncfg_overrides(instance_overrides);
2522        // Re-push the env-wide configs so existing replicas pick up their
2523        // (possibly changed) overrides. This also reverts a removed override:
2524        // the per-replica layer no longer carries the key, so the replica
2525        // falls back to the env-wide value, which both configs always include
2526        // because they render the full dyncfg set.
2527        let compute_config = crate::flags::compute_config(self.catalog().system_config());
2528        self.controller.compute.update_configuration(compute_config);
2529        let storage_config = crate::flags::storage_config(self.catalog().system_config());
2530        self.controller.storage.update_parameters(storage_config);
2531    }
2532
2533    /// Returns the cluster-coherent scoped optimizer-feature overrides for
2534    /// `cluster_id`. See
2535    /// [`CatalogState::cluster_scoped_optimizer_overrides`](crate::catalog::CatalogState::cluster_scoped_optimizer_overrides).
2536    pub(crate) fn cluster_scoped_optimizer_overrides(
2537        &self,
2538        cluster_id: ClusterId,
2539    ) -> OptimizerFeatureOverrides {
2540        self.catalog()
2541            .state()
2542            .cluster_scoped_optimizer_overrides(cluster_id)
2543    }
2544
2545    /// Initializes coordinator state based on the contained catalog. Must be
2546    /// called after creating the coordinator and before calling the
2547    /// `Coordinator::serve` method.
2548    #[instrument(name = "coord::bootstrap")]
2549    pub(crate) async fn bootstrap(
2550        &mut self,
2551        boot_ts: Timestamp,
2552        migrated_storage_collections_0dt: BTreeSet<CatalogItemId>,
2553        hydrate_migrated_mvs: bool,
2554        mut builtin_table_updates: Vec<BuiltinTableUpdate>,
2555        cached_global_exprs: BTreeMap<GlobalId, GlobalExpressions>,
2556        uncached_local_exprs: BTreeMap<GlobalId, LocalExpressions>,
2557    ) -> Result<(), AdapterError> {
2558        let bootstrap_start = Instant::now();
2559        info!("startup: coordinator init: bootstrap beginning");
2560        info!("startup: coordinator init: bootstrap: preamble beginning");
2561
2562        // Initialize cluster replica statuses.
2563        // Gross iterator is to avoid partial borrow issues.
2564        let cluster_statuses: Vec<(_, Vec<_>)> = self
2565            .catalog()
2566            .clusters()
2567            .map(|cluster| {
2568                (
2569                    cluster.id(),
2570                    cluster
2571                        .replicas()
2572                        .map(|replica| {
2573                            (replica.replica_id, replica.config.location.num_processes())
2574                        })
2575                        .collect(),
2576                )
2577            })
2578            .collect();
2579        let now = self.now_datetime();
2580        for (cluster_id, replica_statuses) in cluster_statuses {
2581            self.cluster_replica_statuses
2582                .initialize_cluster_statuses(cluster_id);
2583            for (replica_id, num_processes) in replica_statuses {
2584                self.cluster_replica_statuses
2585                    .initialize_cluster_replica_statuses(
2586                        cluster_id,
2587                        replica_id,
2588                        num_processes,
2589                        now,
2590                    );
2591            }
2592        }
2593
2594        let system_config = self.catalog().system_config();
2595
2596        // Inform metrics about the initial system configuration.
2597        mz_metrics::update_dyncfg(&system_config.dyncfg_updates());
2598
2599        // Inform the controllers about their initial configuration.
2600        let compute_config = flags::compute_config(system_config);
2601        let storage_config = flags::storage_config(system_config);
2602        let scheduling_config = flags::orchestrator_scheduling_config(system_config);
2603        let dyncfg_updates = system_config.dyncfg_updates();
2604        self.controller.compute.update_configuration(compute_config);
2605        self.controller.storage.update_parameters(storage_config);
2606        self.controller
2607            .update_orchestrator_scheduling_config(scheduling_config);
2608        self.controller.update_configuration(dyncfg_updates);
2609
2610        // Install the replica-local scoped overrides before creating any
2611        // replica below. Parts of a replica's configuration (its `TimelyConfig`,
2612        // its expiration offset) are resolved once, when the controller
2613        // provisions the replica, and must see its overrides at that point. The
2614        // push after the creation loop cannot serve this purpose, because those
2615        // values are frozen by then.
2616        let replica_dyncfg_overrides = self.replica_dyncfg_overrides();
2617        self.controller
2618            .update_replica_dyncfg_overrides(replica_dyncfg_overrides);
2619
2620        // Skip the credit consumption check at bootstrap under DisableClusterCreation behavior:
2621        // this codepath validates existing replicas at startup, not cluster creation, so it
2622        // must not block startup. New cluster creation is still gated by the DDL-time check.
2623        // The Disable case is already handled by a bail! in main.rs before we reach here.
2624        let enforce_credit_limit_at_bootstrap = !matches!(
2625            self.license_key.expiration_behavior,
2626            ExpirationBehavior::DisableClusterCreation,
2627        );
2628        if enforce_credit_limit_at_bootstrap {
2629            self.validate_resource_limit_numeric(
2630                Numeric::zero(),
2631                self.current_credit_consumption_rate(None),
2632                |system_vars| {
2633                    self.license_key
2634                        .max_credit_consumption_rate()
2635                        .map_or_else(|| system_vars.max_credit_consumption_rate(), Numeric::from)
2636                },
2637                "cluster replica",
2638                MAX_CREDIT_CONSUMPTION_RATE.name(),
2639            )?;
2640        }
2641
2642        let mut policies_to_set: BTreeMap<CompactionWindow, CollectionIdBundle> =
2643            Default::default();
2644
2645        let enable_worker_core_affinity =
2646            self.catalog().system_config().enable_worker_core_affinity();
2647        for instance in self.catalog.clusters() {
2648            self.controller.create_cluster(
2649                instance.id,
2650                ClusterConfig {
2651                    arranged_logs: instance.log_indexes.clone(),
2652                    workload_class: instance.config.workload_class.clone(),
2653                },
2654            )?;
2655            for replica in instance.replicas() {
2656                let role = instance.role();
2657                self.controller.create_replica(
2658                    instance.id,
2659                    replica.replica_id,
2660                    instance.name.clone(),
2661                    replica.name.clone(),
2662                    role,
2663                    replica.config.clone(),
2664                    enable_worker_core_affinity,
2665                )?;
2666            }
2667        }
2668
2669        // Now that the compute instances and their replicas exist, push the
2670        // replica-local scoped overrides into the controllers so existing
2671        // replicas observe them at startup. The scoped (per-cluster and
2672        // per-replica) working copy was restored from the durable cache into
2673        // `CatalogState` while opening the catalog, so the last-known values are
2674        // in effect before the first parameter sync and through a sync outage.
2675        // This must run after the creation loop above: the push iterates the
2676        // controller's instances, so before they exist it is a no-op. It also
2677        // runs before dataflows are rendered later in bootstrap, so render-frozen
2678        // replica flags take effect. The cluster-coherent layer is read at plan
2679        // time.
2680        self.push_replica_dyncfg_overrides();
2681
2682        info!(
2683            "startup: coordinator init: bootstrap: preamble complete in {:?}",
2684            bootstrap_start.elapsed()
2685        );
2686
2687        let init_storage_collections_start = Instant::now();
2688        info!("startup: coordinator init: bootstrap: storage collections init beginning");
2689        self.bootstrap_storage_collections(&migrated_storage_collections_0dt)
2690            .await;
2691        info!(
2692            "startup: coordinator init: bootstrap: storage collections init complete in {:?}",
2693            init_storage_collections_start.elapsed()
2694        );
2695
2696        // The storage controller knows about the introspection collections now, so we can start
2697        // sinking introspection updates in the compute controller. It makes sense to do that as
2698        // soon as possible, to avoid updates piling up in the compute controller's internal
2699        // buffers.
2700        self.controller.start_compute_introspection_sink();
2701
2702        let sorting_start = Instant::now();
2703        info!("startup: coordinator init: bootstrap: sorting catalog entries");
2704        let entries = self.bootstrap_sort_catalog_entries();
2705        info!(
2706            "startup: coordinator init: bootstrap: sorting catalog entries complete in {:?}",
2707            sorting_start.elapsed()
2708        );
2709
2710        let optimize_dataflows_start = Instant::now();
2711        info!("startup: coordinator init: bootstrap: optimize dataflow plans beginning");
2712        let uncached_global_exps = self.bootstrap_dataflow_plans(&entries, cached_global_exprs)?;
2713        info!(
2714            "startup: coordinator init: bootstrap: optimize dataflow plans complete in {:?}",
2715            optimize_dataflows_start.elapsed()
2716        );
2717
2718        // We don't need to wait for the cache to update.
2719        let _fut = self.catalog().update_expression_cache(
2720            uncached_local_exprs.into_iter().collect(),
2721            uncached_global_exps.into_iter().collect(),
2722            Default::default(),
2723        );
2724
2725        // Select dataflow as-ofs. This step relies on the storage collections created by
2726        // `bootstrap_storage_collections` and the dataflow plans created by
2727        // `bootstrap_dataflow_plans`.
2728        let bootstrap_as_ofs_start = Instant::now();
2729        info!("startup: coordinator init: bootstrap: dataflow as-of bootstrapping beginning");
2730        let dataflow_read_holds = self.bootstrap_dataflow_as_ofs().await;
2731        info!(
2732            "startup: coordinator init: bootstrap: dataflow as-of bootstrapping complete in {:?}",
2733            bootstrap_as_ofs_start.elapsed()
2734        );
2735
2736        let postamble_start = Instant::now();
2737        info!("startup: coordinator init: bootstrap: postamble beginning");
2738
2739        let logs: BTreeSet<_> = BUILTINS::logs()
2740            .map(|log| self.catalog().resolve_builtin_log(log))
2741            .flat_map(|item_id| self.catalog().get_global_ids(&item_id))
2742            .collect();
2743
2744        let mut privatelink_connections = BTreeMap::new();
2745
2746        for entry in &entries {
2747            debug!(
2748                "coordinator init: installing {} {}",
2749                entry.item().typ(),
2750                entry.id()
2751            );
2752            let mut policy = entry.item().initial_logical_compaction_window();
2753            match entry.item() {
2754                // Currently catalog item rebuild assumes that sinks and
2755                // indexes are always built individually and does not store information
2756                // about how it was built. If we start building multiple sinks and/or indexes
2757                // using a single dataflow, we have to make sure the rebuild process re-runs
2758                // the same multiple-build dataflow.
2759                CatalogItem::Source(source) => {
2760                    // Propagate source compaction windows to subsources if needed.
2761                    if source.custom_logical_compaction_window.is_none() {
2762                        if let DataSourceDesc::IngestionExport { ingestion_id, .. } =
2763                            source.data_source
2764                        {
2765                            policy = Some(
2766                                self.catalog()
2767                                    .get_entry(&ingestion_id)
2768                                    .source()
2769                                    .expect("must be source")
2770                                    .custom_logical_compaction_window
2771                                    .unwrap_or_default(),
2772                            );
2773                        }
2774                    }
2775                    policies_to_set
2776                        .entry(policy.expect("sources have a compaction window"))
2777                        .or_insert_with(Default::default)
2778                        .storage_ids
2779                        .insert(source.global_id());
2780                }
2781                CatalogItem::Table(table) => {
2782                    policies_to_set
2783                        .entry(policy.expect("tables have a compaction window"))
2784                        .or_insert_with(Default::default)
2785                        .storage_ids
2786                        .extend(table.global_ids());
2787                }
2788                CatalogItem::Index(idx) => {
2789                    let policy_entry = policies_to_set
2790                        .entry(policy.expect("indexes have a compaction window"))
2791                        .or_insert_with(Default::default);
2792
2793                    if logs.contains(&idx.on) {
2794                        policy_entry
2795                            .compute_ids
2796                            .entry(idx.cluster_id)
2797                            .or_insert_with(BTreeSet::new)
2798                            .insert(idx.global_id());
2799                    } else {
2800                        let df_desc = self
2801                            .catalog()
2802                            .try_get_physical_plan(&idx.global_id())
2803                            .expect("added in `bootstrap_dataflow_plans`")
2804                            .clone();
2805
2806                        let df_meta = self
2807                            .catalog()
2808                            .try_get_dataflow_metainfo(&idx.global_id())
2809                            .expect("added in `bootstrap_dataflow_plans`");
2810
2811                        if self.catalog().state().system_config().enable_mz_notices() {
2812                            // Collect optimization hint updates.
2813                            self.catalog().state().pack_optimizer_notices(
2814                                &mut builtin_table_updates,
2815                                df_meta.optimizer_notices.iter(),
2816                                Diff::ONE,
2817                            );
2818                        }
2819
2820                        // What follows is morally equivalent to `self.ship_dataflow(df, idx.cluster_id)`,
2821                        // but we cannot call that as it will also downgrade the read hold on the index.
2822                        policy_entry
2823                            .compute_ids
2824                            .entry(idx.cluster_id)
2825                            .or_insert_with(Default::default)
2826                            .extend(df_desc.export_ids());
2827
2828                        self.controller
2829                            .compute
2830                            .create_dataflow(idx.cluster_id, df_desc, None)
2831                            .unwrap_or_terminate("cannot fail to create dataflows");
2832                    }
2833                }
2834                CatalogItem::View(_) => (),
2835                CatalogItem::MaterializedView(mview) => {
2836                    // Each version receives a read policy when it is created. Bootstrap
2837                    // must restore every policy because the oldest version owns the shared
2838                    // Persist shard and capability changes reach it through each newer
2839                    // version's primary link. A `NoPolicy` version would block that
2840                    // propagation and pin compaction.
2841                    policies_to_set
2842                        .entry(policy.expect("materialized views have a compaction window"))
2843                        .or_insert_with(Default::default)
2844                        .storage_ids
2845                        .extend(mview.global_ids());
2846
2847                    let mut df_desc = self
2848                        .catalog()
2849                        .try_get_physical_plan(&mview.global_id_writes())
2850                        .expect("added in `bootstrap_dataflow_plans`")
2851                        .clone();
2852
2853                    if let Some(initial_as_of) = mview.initial_as_of.clone() {
2854                        df_desc.set_initial_as_of(initial_as_of);
2855                    }
2856
2857                    // If we have a refresh schedule that has a last refresh, then set the `until` to the last refresh.
2858                    let until = mview
2859                        .refresh_schedule
2860                        .as_ref()
2861                        .and_then(|s| s.last_refresh())
2862                        .and_then(|r| r.try_step_forward());
2863                    if let Some(until) = until {
2864                        df_desc.until.meet_assign(&Antichain::from_elem(until));
2865                    }
2866
2867                    let df_meta = self
2868                        .catalog()
2869                        .try_get_dataflow_metainfo(&mview.global_id_writes())
2870                        .expect("added in `bootstrap_dataflow_plans`");
2871
2872                    if self.catalog().state().system_config().enable_mz_notices() {
2873                        // Collect optimization hint updates.
2874                        self.catalog().state().pack_optimizer_notices(
2875                            &mut builtin_table_updates,
2876                            df_meta.optimizer_notices.iter(),
2877                            Diff::ONE,
2878                        );
2879                    }
2880
2881                    self.ship_dataflow(df_desc, mview.cluster_id, mview.target_replica)
2882                        .await;
2883
2884                    // A pending `REPLACEMENT FOR` MV must stay read-only until
2885                    // `ALTER ... APPLY REPLACEMENT` swaps it in. Unrelated to the
2886                    // builtin-migration `Replacement` mechanism below.
2887                    if mview.replacement_target.is_none() {
2888                        let gid = mview.global_id_writes();
2889                        if hydrate_migrated_mvs
2890                            && migrated_storage_collections_0dt.contains(&entry.id())
2891                        {
2892                            // `migrated_storage_collections_0dt` is `Replacement`-migrated items
2893                            // only, so this is a fresh shard we own: nothing else writes it, and
2894                            // writing it while read-only hydrates the MV and its dependents before
2895                            // cut-over. An `Evolution`-migrated MV reuses the leader's live shard
2896                            // and must never reach here.
2897                            //
2898                            // A *new* builtin MV gets no such treatment: its shard allocation
2899                            // lives only in this read-only savepoint, so the promoted leader
2900                            // allocates a different shard and discards whatever we wrote.
2901                            self.controller
2902                                .compute
2903                                .allow_writes_in_read_only(mview.cluster_id, gid)
2904                                .unwrap_or_terminate("allow_writes cannot fail");
2905                        } else {
2906                            self.allow_writes(mview.cluster_id, gid);
2907                        }
2908                    }
2909                }
2910                CatalogItem::MetricSink(metric_sink) => {
2911                    let df_desc = self
2912                        .catalog()
2913                        .try_get_physical_plan(&metric_sink.global_id)
2914                        .expect("added in `bootstrap_dataflow_plans`")
2915                        .clone();
2916
2917                    let df_meta = self
2918                        .catalog()
2919                        .try_get_dataflow_metainfo(&metric_sink.global_id)
2920                        .expect("added in `bootstrap_dataflow_plans`");
2921
2922                    if self.catalog().state().system_config().enable_mz_notices() {
2923                        // Collect optimization hint updates.
2924                        self.catalog().state().pack_optimizer_notices(
2925                            &mut builtin_table_updates,
2926                            df_meta.optimizer_notices.iter(),
2927                            Diff::ONE,
2928                        );
2929                    }
2930
2931                    // No read policy to set: the export is a sink, not a readable collection, so
2932                    // `ship_dataflow` has no index export to initialize a policy for.
2933                    self.ship_dataflow(df_desc, metric_sink.cluster_id, None)
2934                        .await;
2935                }
2936                CatalogItem::Sink(sink) => {
2937                    policies_to_set
2938                        .entry(CompactionWindow::Default)
2939                        .or_insert_with(Default::default)
2940                        .storage_ids
2941                        .insert(sink.global_id());
2942                }
2943                CatalogItem::Connection(catalog_connection) => {
2944                    if let ConnectionDetails::AwsPrivatelink(conn) = &catalog_connection.details {
2945                        privatelink_connections.insert(
2946                            entry.id(),
2947                            VpcEndpointConfig {
2948                                aws_service_name: conn.service_name.clone(),
2949                                availability_zone_ids: conn.availability_zones.clone(),
2950                            },
2951                        );
2952                    }
2953                }
2954                // Nothing to do for these cases
2955                CatalogItem::Log(_)
2956                | CatalogItem::Type(_)
2957                | CatalogItem::Func(_)
2958                | CatalogItem::Secret(_) => {}
2959            }
2960        }
2961
2962        if let Some(cloud_resource_controller) = &self.cloud_resource_controller {
2963            // Clean up any extraneous VpcEndpoints that shouldn't exist.
2964            let existing_vpc_endpoints = cloud_resource_controller
2965                .list_vpc_endpoints()
2966                .await
2967                .context("list vpc endpoints")?;
2968            let existing_vpc_endpoints = BTreeSet::from_iter(existing_vpc_endpoints.into_keys());
2969            let desired_vpc_endpoints = privatelink_connections.keys().cloned().collect();
2970            let vpc_endpoints_to_remove = existing_vpc_endpoints.difference(&desired_vpc_endpoints);
2971            for id in vpc_endpoints_to_remove {
2972                cloud_resource_controller
2973                    .delete_vpc_endpoint(*id)
2974                    .await
2975                    .context("deleting extraneous vpc endpoint")?;
2976            }
2977
2978            // Ensure desired VpcEndpoints are up to date.
2979            for (id, spec) in privatelink_connections {
2980                cloud_resource_controller
2981                    .ensure_vpc_endpoint(id, spec)
2982                    .await
2983                    .context("ensuring vpc endpoint")?;
2984            }
2985        }
2986
2987        // Having installed all entries, creating all constraints, we can now drop read holds and
2988        // relax read policies.
2989        drop(dataflow_read_holds);
2990        // TODO -- Improve `initialize_read_policies` API so we can avoid calling this in a loop.
2991        for (cw, policies) in policies_to_set {
2992            self.initialize_read_policies(&policies, cw).await;
2993        }
2994
2995        // Expose mapping from T-shirt sizes to actual sizes
2996        builtin_table_updates.extend(
2997            self.catalog().state().resolve_builtin_table_updates(
2998                self.catalog().state().pack_all_replica_size_updates(),
2999            ),
3000        );
3001
3002        debug!("startup: coordinator init: bootstrap: initializing migrated builtin tables");
3003        // When 0dt is enabled, we create new shards for any migrated builtin storage collections.
3004        // In read-only mode, the migrated builtin tables (which are a subset of migrated builtin
3005        // storage collections) need to be back-filled so that any dependent dataflow can be
3006        // hydrated. Additionally, these shards are not registered with the txn-shard, and cannot
3007        // be registered while in read-only, so they are written to directly.
3008        let migrated_updates_fut = if self.controller.read_only() {
3009            let min_timestamp = Timestamp::minimum();
3010            let migrated_builtin_table_updates: Vec<_> = builtin_table_updates
3011                .extract_if(.., |update| {
3012                    let gid = self.catalog().get_entry(&update.id).latest_global_id();
3013                    migrated_storage_collections_0dt.contains(&update.id)
3014                        && self
3015                            .controller
3016                            .storage_collections
3017                            .collection_frontiers(gid)
3018                            .expect("all tables are registered")
3019                            .write_frontier
3020                            .elements()
3021                            == &[min_timestamp]
3022                })
3023                .collect();
3024            if migrated_builtin_table_updates.is_empty() {
3025                futures::future::ready(()).boxed()
3026            } else {
3027                // Group all updates per-table.
3028                let mut grouped_appends: BTreeMap<GlobalId, Vec<TableData>> = BTreeMap::new();
3029                for update in migrated_builtin_table_updates {
3030                    let gid = self.catalog().get_entry(&update.id).latest_global_id();
3031                    grouped_appends.entry(gid).or_default().push(update.data);
3032                }
3033                info!(
3034                    "coordinator init: rehydrating migrated builtin tables in read-only mode: {:?}",
3035                    grouped_appends.keys().collect::<Vec<_>>()
3036                );
3037
3038                // Consolidate Row data, staged batches must already be consolidated.
3039                let mut all_appends = Vec::with_capacity(grouped_appends.len());
3040                for (item_id, table_data) in grouped_appends.into_iter() {
3041                    let mut all_rows = Vec::new();
3042                    let mut all_data = Vec::new();
3043                    for data in table_data {
3044                        match data {
3045                            TableData::Rows(rows) => all_rows.extend(rows),
3046                            TableData::Batches(_) => all_data.push(data),
3047                        }
3048                    }
3049                    differential_dataflow::consolidation::consolidate(&mut all_rows);
3050                    all_data.push(TableData::Rows(all_rows));
3051
3052                    // TODO(parkmycar): Use SmallVec throughout.
3053                    all_appends.push((item_id, all_data));
3054                }
3055
3056                let fut = self
3057                    .controller
3058                    .storage
3059                    .append_table(min_timestamp, boot_ts.step_forward(), all_appends)
3060                    .expect("cannot fail to append");
3061                async {
3062                    fut.await
3063                        .expect("One-shot shouldn't be dropped during bootstrap")
3064                        .unwrap_or_terminate("cannot fail to append")
3065                }
3066                .boxed()
3067            }
3068        } else {
3069            futures::future::ready(()).boxed()
3070        };
3071
3072        info!(
3073            "startup: coordinator init: bootstrap: postamble complete in {:?}",
3074            postamble_start.elapsed()
3075        );
3076
3077        let builtin_update_start = Instant::now();
3078        info!("startup: coordinator init: bootstrap: generate builtin updates beginning");
3079
3080        if self.controller.read_only() {
3081            info!(
3082                "coordinator init: bootstrap: stashing builtin table updates while in read-only mode"
3083            );
3084
3085            self.buffered_builtin_table_updates
3086                .as_mut()
3087                .expect("in read-only mode")
3088                .append(&mut builtin_table_updates);
3089        } else {
3090            self.bootstrap_tables(&entries, builtin_table_updates).await;
3091        };
3092        info!(
3093            "startup: coordinator init: bootstrap: generate builtin updates complete in {:?}",
3094            builtin_update_start.elapsed()
3095        );
3096
3097        let cleanup_secrets_start = Instant::now();
3098        info!("startup: coordinator init: bootstrap: generate secret cleanup beginning");
3099        // Cleanup orphaned secrets. Errors during list() or delete() do not
3100        // need to prevent bootstrap from succeeding; we will retry next
3101        // startup.
3102        {
3103            // Destructure Self so we can selectively move fields into the async
3104            // task.
3105            let Self {
3106                secrets_controller,
3107                catalog,
3108                ..
3109            } = self;
3110
3111            let next_user_item_id = catalog.get_next_user_item_id().await?;
3112            let next_system_item_id = catalog.get_next_system_item_id().await?;
3113            let read_only = self.controller.read_only();
3114            // Fetch all IDs from the catalog to future-proof against other
3115            // things using secrets. Today, SECRET and CONNECTION objects use
3116            // secrets_controller.ensure, but more things could in the future
3117            // that would be easy to miss adding here.
3118            let catalog_ids: BTreeSet<CatalogItemId> =
3119                catalog.entries().map(|entry| entry.id()).collect();
3120            let secrets_controller = Arc::clone(secrets_controller);
3121
3122            spawn(|| "cleanup-orphaned-secrets", async move {
3123                if read_only {
3124                    info!(
3125                        "coordinator init: not cleaning up orphaned secrets while in read-only mode"
3126                    );
3127                    return;
3128                }
3129                info!("coordinator init: cleaning up orphaned secrets");
3130
3131                match secrets_controller.list().await {
3132                    Ok(controller_secrets) => {
3133                        let controller_secrets: BTreeSet<CatalogItemId> =
3134                            controller_secrets.into_iter().collect();
3135                        let orphaned = controller_secrets.difference(&catalog_ids);
3136                        for id in orphaned {
3137                            let id_too_large = match id {
3138                                CatalogItemId::System(id) => *id >= next_system_item_id,
3139                                CatalogItemId::User(id) => *id >= next_user_item_id,
3140                                CatalogItemId::IntrospectionSourceIndex(_)
3141                                | CatalogItemId::Transient(_) => false,
3142                            };
3143                            if id_too_large {
3144                                info!(
3145                                    %next_user_item_id, %next_system_item_id,
3146                                    "coordinator init: not deleting orphaned secret {id} that was likely created by a newer deploy generation"
3147                                );
3148                            } else {
3149                                info!("coordinator init: deleting orphaned secret {id}");
3150                                fail_point!("orphan_secrets");
3151                                if let Err(e) = secrets_controller.delete(*id).await {
3152                                    warn!(
3153                                        "Dropping orphaned secret has encountered an error: {}",
3154                                        e
3155                                    );
3156                                }
3157                            }
3158                        }
3159                    }
3160                    Err(e) => warn!("Failed to list secrets during orphan cleanup: {:?}", e),
3161                }
3162            });
3163        }
3164        info!(
3165            "startup: coordinator init: bootstrap: generate secret cleanup complete in {:?}",
3166            cleanup_secrets_start.elapsed()
3167        );
3168
3169        // Run all of our final steps concurrently.
3170        let final_steps_start = Instant::now();
3171        info!(
3172            "startup: coordinator init: bootstrap: migrate builtin tables in read-only mode beginning"
3173        );
3174        migrated_updates_fut
3175            .instrument(info_span!("coord::bootstrap::final"))
3176            .await;
3177
3178        debug!(
3179            "startup: coordinator init: bootstrap: announcing completion of initialization to controller"
3180        );
3181        // Announce the completion of initialization.
3182        self.controller.initialization_complete();
3183
3184        // Initialize unified introspection.
3185        self.bootstrap_introspection_subscribes().await;
3186
3187        // Install the curated metric sinks on every replica.
3188        self.bootstrap_metric_sinks().await;
3189
3190        info!(
3191            "startup: coordinator init: bootstrap: migrate builtin tables in read-only mode complete in {:?}",
3192            final_steps_start.elapsed()
3193        );
3194
3195        info!(
3196            "startup: coordinator init: bootstrap complete in {:?}",
3197            bootstrap_start.elapsed()
3198        );
3199        Ok(())
3200    }
3201
3202    /// Prepares tables for writing by resetting them to a known state and
3203    /// appending the given builtin table updates. The timestamp oracle
3204    /// will be advanced to the write timestamp of the append when this
3205    /// method returns.
3206    #[allow(clippy::async_yields_async)]
3207    #[instrument]
3208    async fn bootstrap_tables(
3209        &mut self,
3210        entries: &[CatalogEntry],
3211        mut builtin_table_updates: Vec<BuiltinTableUpdate>,
3212    ) {
3213        /// Smaller helper struct of metadata for bootstrapping tables.
3214        struct TableMetadata<'a> {
3215            id: CatalogItemId,
3216            name: &'a QualifiedItemName,
3217            table: &'a Table,
3218        }
3219
3220        // Filter our entries down to just tables.
3221        let table_metas: Vec<_> = entries
3222            .into_iter()
3223            .filter_map(|entry| {
3224                entry.table().map(|table| TableMetadata {
3225                    id: entry.id(),
3226                    name: entry.name(),
3227                    table,
3228                })
3229            })
3230            .collect();
3231
3232        // Append empty batches to advance the timestamp of all tables.
3233        debug!("coordinator init: advancing all tables to current timestamp");
3234        let WriteTimestamp {
3235            timestamp: write_ts,
3236            advance_to,
3237        } = self.get_local_write_ts().await;
3238        let appends = table_metas
3239            .iter()
3240            .map(|meta| (meta.table.global_id_writes(), Vec::new()))
3241            .collect();
3242        // Append the tables in the background. We apply the write timestamp before getting a read
3243        // timestamp and reading a snapshot of each table, so the snapshots will block on their own
3244        // until the appends are complete.
3245        let table_fence_rx = self
3246            .controller
3247            .storage
3248            .append_table(write_ts.clone(), advance_to, appends)
3249            .expect("invalid updates");
3250
3251        self.apply_local_write(write_ts).await;
3252
3253        // Add builtin table updates the clear the contents of all system tables
3254        debug!("coordinator init: resetting system tables");
3255        let read_ts = self.get_local_read_ts().await;
3256
3257        let retained_across_restarts = BTreeSet::from([
3258            self.catalog()
3259                .resolve_builtin_table(&MZ_STORAGE_USAGE_BY_SHARD),
3260            self.catalog()
3261                .resolve_builtin_table(&MZ_OBJECT_ARRANGEMENT_SIZE_HISTORY),
3262            self.catalog()
3263                .resolve_builtin_table(&MZ_OBJECT_HYDRATION_HISTORY),
3264            self.catalog()
3265                .resolve_builtin_table(&MZ_REPLICA_HYDRATION_HISTORY),
3266        ]);
3267
3268        let mut retraction_tasks = Vec::new();
3269        let system_tables: Vec<_> = table_metas
3270            .iter()
3271            .filter(|meta| meta.id.is_system() && !retained_across_restarts.contains(&meta.id))
3272            .collect();
3273
3274        for system_table in system_tables {
3275            let table_id = system_table.id;
3276            let full_name = self.catalog().resolve_full_name(system_table.name, None);
3277            debug!("coordinator init: resetting system table {full_name} ({table_id})");
3278
3279            // Fetch the current contents of the table for retraction.
3280            let snapshot_fut = self
3281                .controller
3282                .storage_collections
3283                .snapshot_cursor(system_table.table.global_id_writes(), read_ts);
3284            let batch_fut = self
3285                .controller
3286                .storage_collections
3287                .create_update_builder(system_table.table.global_id_writes());
3288
3289            let task = spawn(|| format!("snapshot-{table_id}"), async move {
3290                // Create a TimestamplessUpdateBuilder.
3291                let mut batch = batch_fut
3292                    .await
3293                    .unwrap_or_terminate("cannot fail to create a batch for a BuiltinTable");
3294                tracing::info!(?table_id, "starting snapshot");
3295                // Get a cursor which will emit a consolidated snapshot.
3296                let mut snapshot_cursor = snapshot_fut
3297                    .await
3298                    .unwrap_or_terminate("cannot fail to snapshot");
3299
3300                // Retract the current contents, spilling into our builder.
3301                while let Some(values) = snapshot_cursor.next().await {
3302                    for (key, _t, d) in values {
3303                        let d_invert = d.neg();
3304                        batch.add(&key, &(), &d_invert).await;
3305                    }
3306                }
3307                tracing::info!(?table_id, "finished snapshot");
3308
3309                let batch = batch.finish().await;
3310                BuiltinTableUpdate::batch(table_id, batch)
3311            });
3312            retraction_tasks.push(task);
3313        }
3314
3315        let retractions_res = futures::future::join_all(retraction_tasks).await;
3316        for retractions in retractions_res {
3317            builtin_table_updates.push(retractions);
3318        }
3319
3320        // Now that the snapshots are complete, the appends must also be complete.
3321        table_fence_rx
3322            .await
3323            .expect("One-shot shouldn't be dropped during bootstrap")
3324            .unwrap_or_terminate("cannot fail to append");
3325
3326        info!("coordinator init: sending builtin table updates");
3327        let builtin_updates_fut = self.builtin_table_update().execute(builtin_table_updates);
3328        // Wait for the committer to apply the write, so the builtin tables are readable before
3329        // we start serving. The committer allocates the timestamp and advances the oracle.
3330        builtin_updates_fut.await;
3331    }
3332
3333    /// Initializes all storage collections required by catalog objects in the storage controller.
3334    ///
3335    /// This method takes care of collection creation, as well as migration of existing
3336    /// collections.
3337    ///
3338    /// Creating all storage collections in a single `create_collections` call, rather than on
3339    /// demand, is more efficient as it reduces the number of writes to durable storage. It also
3340    /// allows subsequent bootstrap logic to fetch metadata (such as frontiers) of arbitrary
3341    /// storage collections, without needing to worry about dependency order.
3342    ///
3343    /// `migrated_storage_collections` is a set of builtin storage collections that have been
3344    /// migrated and should be handled specially.
3345    #[instrument]
3346    async fn bootstrap_storage_collections(
3347        &mut self,
3348        migrated_storage_collections: &BTreeSet<CatalogItemId>,
3349    ) {
3350        let catalog = self.catalog();
3351
3352        let source_desc = |object_id: GlobalId,
3353                           data_source: &DataSourceDesc,
3354                           desc: &RelationDesc,
3355                           timeline: &Timeline| {
3356            let data_source = match data_source.clone() {
3357                // Re-announce the source description.
3358                DataSourceDesc::Ingestion { desc, cluster_id } => {
3359                    let desc = desc.into_inline_connection(catalog.state());
3360                    let ingestion = IngestionDescription::new(desc, cluster_id, object_id);
3361                    DataSource::Ingestion(ingestion)
3362                }
3363                DataSourceDesc::OldSyntaxIngestion {
3364                    desc,
3365                    progress_subsource,
3366                    data_config,
3367                    details,
3368                    cluster_id,
3369                } => {
3370                    let desc = desc.into_inline_connection(catalog.state());
3371                    let data_config = data_config.into_inline_connection(catalog.state());
3372                    // TODO(parkmycar): We should probably check the type here, but I'm not sure if
3373                    // this will always be a Source or a Table.
3374                    let progress_subsource =
3375                        catalog.get_entry(&progress_subsource).latest_global_id();
3376                    let mut ingestion =
3377                        IngestionDescription::new(desc, cluster_id, progress_subsource);
3378                    let legacy_export = SourceExport {
3379                        storage_metadata: (),
3380                        data_config,
3381                        details,
3382                    };
3383                    ingestion.source_exports.insert(object_id, legacy_export);
3384
3385                    DataSource::Ingestion(ingestion)
3386                }
3387                DataSourceDesc::IngestionExport {
3388                    ingestion_id,
3389                    external_reference: _,
3390                    details,
3391                    data_config,
3392                } => {
3393                    // TODO(parkmycar): We should probably check the type here, but I'm not sure if
3394                    // this will always be a Source or a Table.
3395                    let ingestion_id = catalog.get_entry(&ingestion_id).latest_global_id();
3396
3397                    DataSource::IngestionExport {
3398                        ingestion_id,
3399                        details,
3400                        data_config: data_config.into_inline_connection(catalog.state()),
3401                    }
3402                }
3403                DataSourceDesc::Webhook { .. } => DataSource::Webhook,
3404                DataSourceDesc::Progress => DataSource::Progress,
3405                DataSourceDesc::Introspection(introspection) => {
3406                    DataSource::Introspection(introspection)
3407                }
3408                DataSourceDesc::Catalog => DataSource::Other,
3409            };
3410            CollectionDescription {
3411                desc: desc.clone(),
3412                data_source,
3413                since: None,
3414                timeline: Some(timeline.clone()),
3415                primary: None,
3416            }
3417        };
3418
3419        let mut compute_collections = vec![];
3420        let mut collections = vec![];
3421        for entry in catalog.entries() {
3422            match entry.item() {
3423                CatalogItem::Source(source) => {
3424                    collections.push((
3425                        source.global_id(),
3426                        source_desc(
3427                            source.global_id(),
3428                            &source.data_source,
3429                            &source.desc,
3430                            &source.timeline,
3431                        ),
3432                    ));
3433                }
3434                CatalogItem::Table(table) => {
3435                    match &table.data_source {
3436                        TableDataSource::TableWrites { defaults: _ } => {
3437                            let versions: BTreeMap<_, _> = table
3438                                .collection_descs()
3439                                .map(|(gid, version, desc)| (version, (gid, desc)))
3440                                .collect();
3441                            let collection_descs = versions.iter().map(|(version, (gid, desc))| {
3442                                let next_version = version.bump();
3443                                let primary_collection =
3444                                    versions.get(&next_version).map(|(gid, _desc)| gid).copied();
3445                                let mut collection_desc =
3446                                    CollectionDescription::for_table(desc.clone());
3447                                collection_desc.primary = primary_collection;
3448
3449                                (*gid, collection_desc)
3450                            });
3451                            collections.extend(collection_descs);
3452                        }
3453                        TableDataSource::DataSource {
3454                            desc: data_source_desc,
3455                            timeline,
3456                        } => {
3457                            // TODO(alter_table): Support versioning tables that read from sources.
3458                            soft_assert_eq_or_log!(table.collections.len(), 1);
3459                            let collection_descs =
3460                                table.collection_descs().map(|(gid, _version, desc)| {
3461                                    (
3462                                        gid,
3463                                        source_desc(
3464                                            entry.latest_global_id(),
3465                                            data_source_desc,
3466                                            &desc,
3467                                            timeline,
3468                                        ),
3469                                    )
3470                                });
3471                            collections.extend(collection_descs);
3472                        }
3473                    };
3474                }
3475                CatalogItem::MaterializedView(mv) => {
3476                    // Applying a replacement preserves the ownership link established when the
3477                    // replacement was created. The oldest collection owns the shard, each applied
3478                    // replacement points to its predecessor, and a pending replacement starts by
3479                    // pointing to its target's latest collection.
3480                    //
3481                    // NOTE: Versioned tables chain in the opposite direction because their latest
3482                    // version owns the shard. Each chain matches its runtime replacement path.
3483                    let mut primary = mv
3484                        .replacement_target
3485                        .map(|target_id| catalog.get_entry(&target_id).latest_global_id());
3486                    let collection_descs = mv.collection_descs().map(|(gid, _version, desc)| {
3487                        let mut collection_desc =
3488                            CollectionDescription::for_other(desc, mv.initial_as_of.clone());
3489                        collection_desc.primary = primary;
3490                        primary = Some(gid);
3491                        (gid, collection_desc)
3492                    });
3493
3494                    collections.extend(collection_descs);
3495                    compute_collections.push((mv.global_id_writes(), mv.desc.latest()));
3496                }
3497                CatalogItem::Sink(sink) => {
3498                    let storage_sink_from_entry = self.catalog().get_entry_by_global_id(&sink.from);
3499                    let from_desc = storage_sink_from_entry
3500                        .relation_desc()
3501                        .expect("sinks can only be built on items with descs")
3502                        .into_owned();
3503                    let collection_desc = CollectionDescription {
3504                        // TODO(sinks): make generic once we have more than one sink type.
3505                        desc: KAFKA_PROGRESS_DESC.clone(),
3506                        data_source: DataSource::Sink {
3507                            desc: ExportDescription {
3508                                sink: StorageSinkDesc {
3509                                    from: sink.from,
3510                                    from_desc,
3511                                    connection: sink
3512                                        .connection
3513                                        .clone()
3514                                        .into_inline_connection(self.catalog().state()),
3515                                    envelope: sink.envelope,
3516                                    as_of: Antichain::from_elem(Timestamp::minimum()),
3517                                    with_snapshot: sink.with_snapshot,
3518                                    version: sink.version,
3519                                    from_storage_metadata: (),
3520                                    to_storage_metadata: (),
3521                                    commit_interval: sink.commit_interval,
3522                                },
3523                                instance_id: sink.cluster_id,
3524                            },
3525                        },
3526                        since: None,
3527                        timeline: None,
3528                        primary: None,
3529                    };
3530                    collections.push((sink.global_id, collection_desc));
3531                }
3532                CatalogItem::Log(_)
3533                | CatalogItem::View(_)
3534                | CatalogItem::Index(_)
3535                | CatalogItem::Type(_)
3536                | CatalogItem::Func(_)
3537                | CatalogItem::Secret(_)
3538                | CatalogItem::Connection(_)
3539                // Nothing to bootstrap: a metric sink has no storage collection, it publishes
3540                // into the replica's metrics registry.
3541                | CatalogItem::MetricSink(_) => (),
3542            }
3543        }
3544
3545        let register_ts = if self.controller.read_only() {
3546            self.get_local_read_ts().await
3547        } else {
3548            // Getting a write timestamp bumps the write timestamp in the
3549            // oracle, which we're not allowed in read-only mode.
3550            self.get_local_write_ts().await.timestamp
3551        };
3552
3553        let storage_metadata = self.catalog.state().storage_metadata();
3554        let migrated_storage_collections = migrated_storage_collections
3555            .into_iter()
3556            .flat_map(|item_id| self.catalog.get_entry(item_id).global_ids())
3557            .collect();
3558
3559        // Before possibly creating collections, make sure their schemas are correct.
3560        //
3561        // Across different versions of Materialize the nullability of columns can change based on
3562        // updates to our optimizer.
3563        self.controller
3564            .storage
3565            .evolve_nullability_for_bootstrap(storage_metadata, compute_collections)
3566            .await
3567            .unwrap_or_terminate("cannot fail to evolve collections");
3568
3569        // New builtin storage collections are by default created with [0] since/upper frontiers.
3570        // For collections that have dependencies on other collections (MVs, CTs), this can violate
3571        // the frontier invariants assumed by as-of selection. For example, as-of selection expects
3572        // to be able to pick up computing a materialized view from its most recent upper, but if
3573        // that upper is [0] it's likely that the required times are not available anymore in the
3574        // MV inputs.
3575        //
3576        // To avoid violating frontier invariants, we need to bump their sinces to times greater
3577        // than all of their upstream storage inputs. To know the since of a storage input, it has
3578        // to be registered with the storage controller first. Thus we register collections in
3579        // layers: Each iteration registers the collections whose dependencies are all already
3580        // registered.
3581        let mut pending: BTreeMap<_, _> = collections.into_iter().collect();
3582
3583        // Precompute storage-collection dependencies for each collection.
3584        let transitive_dep_gids: BTreeMap<_, _> = pending
3585            .keys()
3586            .map(|gid| {
3587                let entry = self.catalog.get_entry_by_global_id(gid);
3588                let item_id = entry.id();
3589                let deps = self.catalog.state().transitive_uses(item_id);
3590                let dep_gids: BTreeSet<_> = deps
3591                    // Ignore self-dependencies. For example, `transitive_uses` includes the input ID,
3592                    // and CTs can depend on themselves.
3593                    .filter(|dep_id| *dep_id != item_id)
3594                    .map(|dep_id| self.catalog.get_entry(&dep_id).latest_global_id())
3595                    // Ignore dependencies on objects that are not storage collections.
3596                    .filter(|dep_gid| pending.contains_key(dep_gid))
3597                    .collect();
3598                (*gid, dep_gids)
3599            })
3600            .collect();
3601
3602        let mut created_gids = Vec::new();
3603
3604        while !pending.is_empty() {
3605            // Drain collections whose dependencies have all been registered already
3606            // (i.e., are not in `pending`).
3607            let ready_gids: BTreeSet<_> = pending
3608                .keys()
3609                .filter(|gid| {
3610                    let mut deps = transitive_dep_gids[gid].iter();
3611                    !deps.any(|dep_gid| pending.contains_key(dep_gid))
3612                })
3613                .copied()
3614                .collect();
3615            let mut ready: Vec<_> = pending
3616                .extract_if(.., |gid, _| ready_gids.contains(gid))
3617                .collect();
3618
3619            // Bump sinces of builtin collections.
3620            for (gid, collection) in &mut ready {
3621                // Don't silently overwrite an explicitly specified `since`.
3622                if !gid.is_system() || collection.since.is_some() {
3623                    continue;
3624                }
3625
3626                let mut derived_since = Antichain::from_elem(Timestamp::MIN);
3627                for dep_gid in &transitive_dep_gids[gid] {
3628                    let (since, _) = self
3629                        .controller
3630                        .storage
3631                        .collection_frontiers(*dep_gid)
3632                        .expect("previously registered");
3633                    derived_since.join_assign(&since);
3634                }
3635                collection.since = Some(derived_since);
3636            }
3637
3638            if ready.is_empty() {
3639                soft_panic_or_log!(
3640                    "cycle in storage collections: {:?}",
3641                    pending.keys().collect::<Vec<_>>(),
3642                );
3643                // We get here only due to a bug. Rather than crash-looping, we try our best to
3644                // reach a sane state by attempting to register all the remaining collections at
3645                // once.
3646                ready = mem::take(&mut pending).into_iter().collect();
3647            }
3648
3649            created_gids.extend(ready.iter().map(|(gid, _collection)| *gid));
3650
3651            self.controller
3652                .storage
3653                .create_collections_for_bootstrap(
3654                    storage_metadata,
3655                    Some(register_ts),
3656                    ready,
3657                    &migrated_storage_collections,
3658                )
3659                .await
3660                .unwrap_or_terminate("cannot fail to create collections");
3661        }
3662
3663        // Register txn-wal tables before the later system-table snapshot.
3664        self.controller
3665            .storage
3666            .register_table_collections(register_ts, created_gids)
3667            .await
3668            .unwrap_or_terminate("cannot fail to register tables");
3669
3670        if !self.controller.read_only() {
3671            self.apply_local_write(register_ts).await;
3672        }
3673    }
3674
3675    /// Returns the current list of catalog entries, sorted into an appropriate order for
3676    /// bootstrapping.
3677    ///
3678    /// The returned entries are in dependency order. Indexes are sorted immediately after the
3679    /// objects they index, to ensure that all dependants of these indexed objects can make use of
3680    /// the respective indexes.
3681    fn bootstrap_sort_catalog_entries(&self) -> Vec<CatalogEntry> {
3682        let mut indexes_on = BTreeMap::<_, Vec<_>>::new();
3683        let mut non_indexes = Vec::new();
3684        for entry in self.catalog().entries().cloned() {
3685            if let Some(index) = entry.index() {
3686                let on = self.catalog().get_entry_by_global_id(&index.on);
3687                indexes_on.entry(on.id()).or_default().push(entry);
3688            } else {
3689                non_indexes.push(entry);
3690            }
3691        }
3692
3693        let key_fn = |entry: &CatalogEntry| entry.id;
3694        let dependencies_fn = |entry: &CatalogEntry| entry.uses();
3695        sort_topological(&mut non_indexes, key_fn, dependencies_fn);
3696
3697        let mut result = Vec::new();
3698        for entry in non_indexes {
3699            let id = entry.id();
3700            result.push(entry);
3701            if let Some(mut indexes) = indexes_on.remove(&id) {
3702                result.append(&mut indexes);
3703            }
3704        }
3705
3706        soft_assert_or_log!(
3707            indexes_on.is_empty(),
3708            "indexes with missing dependencies: {indexes_on:?}",
3709        );
3710
3711        result
3712    }
3713
3714    /// Invokes the optimizer on all indexes and materialized views in the catalog and inserts the
3715    /// resulting dataflow plans into the catalog state.
3716    ///
3717    /// `ordered_catalog_entries` must be sorted in dependency order, with dependencies ordered
3718    /// before their dependants.
3719    ///
3720    /// This method does not perform timestamp selection for the dataflows, nor does it create them
3721    /// in the compute controller. Both of these steps happen later during bootstrapping.
3722    ///
3723    /// Returns a map of expressions that were not cached.
3724    #[instrument]
3725    fn bootstrap_dataflow_plans(
3726        &mut self,
3727        ordered_catalog_entries: &[CatalogEntry],
3728        mut cached_global_exprs: BTreeMap<GlobalId, GlobalExpressions>,
3729    ) -> Result<BTreeMap<GlobalId, GlobalExpressions>, AdapterError> {
3730        // The optimizer expects to be able to query its `ComputeInstanceSnapshot` for
3731        // collections the current dataflow can depend on. But since we don't yet install anything
3732        // on compute instances, the snapshot information is incomplete. We fix that by manually
3733        // updating `ComputeInstanceSnapshot` objects to ensure they contain collections previously
3734        // optimized.
3735        let mut instance_snapshots = BTreeMap::new();
3736        let mut uncached_expressions = BTreeMap::new();
3737
3738        let optimizer_config = |catalog: &Catalog, cluster_id| {
3739            let system_config = catalog.system_config();
3740            let overrides = catalog.get_cluster(cluster_id).config.features();
3741            OptimizerConfig::from(system_config)
3742                .override_from(&overrides)
3743                // A cluster-scoped LaunchDarkly rule beats a manual `FEATURES`
3744                // pin.
3745                .override_from(
3746                    &catalog
3747                        .state()
3748                        .cluster_scoped_optimizer_overrides(cluster_id),
3749                )
3750        };
3751
3752        for entry in ordered_catalog_entries {
3753            match entry.item() {
3754                CatalogItem::Index(idx) => {
3755                    // Collect optimizer parameters.
3756                    let compute_instance =
3757                        instance_snapshots.entry(idx.cluster_id).or_insert_with(|| {
3758                            self.instance_snapshot(idx.cluster_id)
3759                                .expect("compute instance exists")
3760                        });
3761                    let global_id = idx.global_id();
3762
3763                    // The index may already be installed on the compute instance. For example,
3764                    // this is the case for introspection indexes.
3765                    if compute_instance.contains_collection(&global_id) {
3766                        continue;
3767                    }
3768
3769                    let optimizer_config = optimizer_config(&self.catalog, idx.cluster_id);
3770
3771                    let (optimized_plan, physical_plan, metainfo) =
3772                        match cached_global_exprs.remove(&global_id) {
3773                            Some(global_expressions)
3774                                if global_expressions.optimizer_features
3775                                    == optimizer_config.features =>
3776                            {
3777                                debug!("global expression cache hit for {global_id:?}");
3778                                (
3779                                    global_expressions.global_mir,
3780                                    global_expressions.physical_plan,
3781                                    global_expressions.dataflow_metainfos,
3782                                )
3783                            }
3784                            Some(_) | None => {
3785                                let (optimized_plan, global_lir_plan) = {
3786                                    // Build an optimizer for this INDEX.
3787                                    let mut optimizer = optimize::index::Optimizer::new(
3788                                        self.owned_catalog(),
3789                                        compute_instance.clone(),
3790                                        global_id,
3791                                        optimizer_config.clone(),
3792                                        self.optimizer_metrics(),
3793                                    );
3794
3795                                    // MIR ⇒ MIR optimization (global)
3796                                    let index_plan = optimize::index::Index::new(
3797                                        entry.name().clone(),
3798                                        idx.on,
3799                                        idx.keys.to_vec(),
3800                                    );
3801                                    let global_mir_plan = optimizer.optimize(index_plan)?;
3802                                    let optimized_plan = global_mir_plan.df_desc().clone();
3803
3804                                    // MIR ⇒ LIR lowering and LIR ⇒ LIR optimization (global)
3805                                    let global_lir_plan = optimizer.optimize(global_mir_plan)?;
3806
3807                                    (optimized_plan, global_lir_plan)
3808                                };
3809
3810                                let (physical_plan, metainfo) = global_lir_plan.unapply();
3811                                let metainfo = {
3812                                    // Pre-allocate a vector of transient GlobalIds for each notice.
3813                                    let notice_ids =
3814                                        std::iter::repeat_with(|| self.allocate_transient_id())
3815                                            .map(|(_item_id, gid)| gid)
3816                                            .take(metainfo.optimizer_notices.len())
3817                                            .collect::<Vec<_>>();
3818                                    // Return a metainfo with rendered notices.
3819                                    self.catalog().render_notices(
3820                                        metainfo,
3821                                        notice_ids,
3822                                        Some(idx.global_id()),
3823                                    )
3824                                };
3825                                uncached_expressions.insert(
3826                                    global_id,
3827                                    GlobalExpressions {
3828                                        global_mir: optimized_plan.clone(),
3829                                        physical_plan: physical_plan.clone(),
3830                                        dataflow_metainfos: metainfo.clone(),
3831                                        optimizer_features: optimizer_config.features.clone(),
3832                                        item_version: RelationVersion::root(),
3833                                    },
3834                                );
3835                                (optimized_plan, physical_plan, metainfo)
3836                            }
3837                        };
3838
3839                    let catalog = self.catalog_mut();
3840                    catalog.set_optimized_plan(idx.global_id(), optimized_plan);
3841                    catalog.set_physical_plan(idx.global_id(), physical_plan);
3842                    catalog.set_dataflow_metainfo(idx.global_id(), metainfo);
3843
3844                    compute_instance.insert_collection(idx.global_id());
3845                }
3846                CatalogItem::MaterializedView(mv) => {
3847                    // Collect optimizer parameters.
3848                    let compute_instance =
3849                        instance_snapshots.entry(mv.cluster_id).or_insert_with(|| {
3850                            self.instance_snapshot(mv.cluster_id)
3851                                .expect("compute instance exists")
3852                        });
3853                    let global_id = mv.global_id_writes();
3854
3855                    let optimizer_config = optimizer_config(&self.catalog, mv.cluster_id);
3856
3857                    let (optimized_plan, physical_plan, metainfo) = match cached_global_exprs
3858                        .remove(&global_id)
3859                    {
3860                        Some(global_expressions)
3861                            if global_expressions.optimizer_features
3862                                == optimizer_config.features =>
3863                        {
3864                            debug!("global expression cache hit for {global_id:?}");
3865                            (
3866                                global_expressions.global_mir,
3867                                global_expressions.physical_plan,
3868                                global_expressions.dataflow_metainfos,
3869                            )
3870                        }
3871                        Some(_) | None => {
3872                            let (_, internal_view_id) = self.allocate_transient_id();
3873                            let debug_name = self
3874                                .catalog()
3875                                .resolve_full_name(entry.name(), None)
3876                                .to_string();
3877
3878                            let (optimized_plan, global_lir_plan) = {
3879                                // Build an optimizer for this MATERIALIZED VIEW.
3880                                let mut optimizer = optimize::materialized_view::Optimizer::new(
3881                                    self.owned_catalog().as_optimizer_catalog(),
3882                                    compute_instance.clone(),
3883                                    global_id,
3884                                    internal_view_id,
3885                                    mv.desc.latest().iter_names().cloned().collect(),
3886                                    mv.non_null_assertions.clone(),
3887                                    mv.refresh_schedule.clone(),
3888                                    debug_name,
3889                                    optimizer_config.clone(),
3890                                    self.optimizer_metrics(),
3891                                );
3892
3893                                // MIR ⇒ MIR optimization (global)
3894                                // We make sure to use the HIR SQL type (since MIR SQL types may not be coherent).
3895                                let typ = infer_sql_type_for_catalog(
3896                                    &mv.raw_expr,
3897                                    &mv.locally_optimized_expr.as_ref().clone(),
3898                                );
3899                                let global_mir_plan = optimizer
3900                                    .optimize((mv.locally_optimized_expr.as_ref().clone(), typ))?;
3901                                let optimized_plan = global_mir_plan.df_desc().clone();
3902
3903                                // MIR ⇒ LIR lowering and LIR ⇒ LIR optimization (global)
3904                                let global_lir_plan = optimizer.optimize(global_mir_plan)?;
3905
3906                                (optimized_plan, global_lir_plan)
3907                            };
3908
3909                            let (physical_plan, metainfo) = global_lir_plan.unapply();
3910                            let metainfo = {
3911                                // Pre-allocate a vector of transient GlobalIds for each notice.
3912                                let notice_ids =
3913                                    std::iter::repeat_with(|| self.allocate_transient_id())
3914                                        .map(|(_item_id, global_id)| global_id)
3915                                        .take(metainfo.optimizer_notices.len())
3916                                        .collect::<Vec<_>>();
3917                                // Return a metainfo with rendered notices.
3918                                self.catalog().render_notices(
3919                                    metainfo,
3920                                    notice_ids,
3921                                    Some(mv.global_id_writes()),
3922                                )
3923                            };
3924                            uncached_expressions.insert(
3925                                global_id,
3926                                GlobalExpressions {
3927                                    global_mir: optimized_plan.clone(),
3928                                    physical_plan: physical_plan.clone(),
3929                                    dataflow_metainfos: metainfo.clone(),
3930                                    optimizer_features: optimizer_config.features.clone(),
3931                                    item_version: latest_item_version(&mv.collections),
3932                                },
3933                            );
3934                            (optimized_plan, physical_plan, metainfo)
3935                        }
3936                    };
3937
3938                    let catalog = self.catalog_mut();
3939                    catalog.set_optimized_plan(mv.global_id_writes(), optimized_plan);
3940                    catalog.set_physical_plan(mv.global_id_writes(), physical_plan);
3941                    catalog.set_dataflow_metainfo(mv.global_id_writes(), metainfo);
3942
3943                    compute_instance.insert_collection(mv.global_id_writes());
3944                }
3945                CatalogItem::MetricSink(metric_sink) => {
3946                    // Collect optimizer parameters.
3947                    let compute_instance = instance_snapshots
3948                        .entry(metric_sink.cluster_id)
3949                        .or_insert_with(|| {
3950                            self.instance_snapshot(metric_sink.cluster_id)
3951                                .expect("compute instance exists")
3952                        });
3953                    let global_id = metric_sink.global_id;
3954                    let optimizer_config = optimizer_config(&self.catalog, metric_sink.cluster_id);
3955
3956                    let (optimized_plan, physical_plan, metainfo) = match cached_global_exprs
3957                        .remove(&global_id)
3958                    {
3959                        Some(global_expressions)
3960                            if global_expressions.optimizer_features
3961                                == optimizer_config.features =>
3962                        {
3963                            debug!("global expression cache hit for {global_id:?}");
3964                            (
3965                                global_expressions.global_mir,
3966                                global_expressions.physical_plan,
3967                                global_expressions.dataflow_metainfos,
3968                            )
3969                        }
3970                        Some(_) | None => {
3971                            // A transient id for the view the optimizer builds over `from` to
3972                            // shape its rows (see `optimize::metric_sink::shape_metric_sink_source`).
3973                            // The id only needs to be unique within this dataflow, so a cached plan
3974                            // reusing a transient id from a previous boot is safe: build ids are
3975                            // dataflow-local on the worker and never registered in the controller's
3976                            // instance-global collections (only export ids are).
3977                            let (_, view_id) = self.allocate_transient_id();
3978
3979                            let (optimized_plan, global_lir_plan) = {
3980                                let mut optimizer = optimize::metric_sink::Optimizer::new(
3981                                    self.owned_catalog(),
3982                                    compute_instance.clone(),
3983                                    view_id,
3984                                    global_id,
3985                                    optimizer_config.clone(),
3986                                    self.optimizer_metrics(),
3987                                );
3988
3989                                // MIR ⇒ MIR optimization (global)
3990                                let metric_sink_plan = optimize::metric_sink::MetricSink::new(
3991                                    self.catalog()
3992                                        .resolve_full_name(entry.name(), None)
3993                                        .to_string(),
3994                                    optimize::metric_sink::MetricSinkFrom::Id(metric_sink.from),
3995                                    metric_sink.prefix.clone(),
3996                                    None,
3997                                );
3998                                let global_mir_plan = optimizer.optimize(metric_sink_plan)?;
3999                                let optimized_plan = global_mir_plan.df_desc().clone();
4000
4001                                // MIR ⇒ LIR lowering and LIR ⇒ LIR optimization (global)
4002                                let global_lir_plan = optimizer.optimize(global_mir_plan)?;
4003
4004                                (optimized_plan, global_lir_plan)
4005                            };
4006
4007                            let (physical_plan, metainfo) = global_lir_plan.unapply();
4008                            let metainfo = {
4009                                // Pre-allocate a vector of transient GlobalIds for each notice.
4010                                let notice_ids =
4011                                    std::iter::repeat_with(|| self.allocate_transient_id())
4012                                        .map(|(_item_id, gid)| gid)
4013                                        .take(metainfo.optimizer_notices.len())
4014                                        .collect::<Vec<_>>();
4015                                // Return a metainfo with rendered notices.
4016                                self.catalog()
4017                                    .render_notices(metainfo, notice_ids, Some(global_id))
4018                            };
4019                            uncached_expressions.insert(
4020                                global_id,
4021                                GlobalExpressions {
4022                                    global_mir: optimized_plan.clone(),
4023                                    physical_plan: physical_plan.clone(),
4024                                    dataflow_metainfos: metainfo.clone(),
4025                                    optimizer_features: optimizer_config.features.clone(),
4026                                    item_version: RelationVersion::root(),
4027                                },
4028                            );
4029                            (optimized_plan, physical_plan, metainfo)
4030                        }
4031                    };
4032
4033                    let catalog = self.catalog_mut();
4034                    catalog.set_optimized_plan(global_id, optimized_plan);
4035                    catalog.set_physical_plan(global_id, physical_plan);
4036                    catalog.set_dataflow_metainfo(global_id, metainfo);
4037
4038                    // NOTE: No `insert_collection` for the export. A metric sink writes to the
4039                    // metrics registry rather than to a readable collection, so no later dataflow
4040                    // can import it.
4041                }
4042                CatalogItem::Table(_)
4043                | CatalogItem::Source(_)
4044                | CatalogItem::Log(_)
4045                | CatalogItem::View(_)
4046                | CatalogItem::Sink(_)
4047                | CatalogItem::Type(_)
4048                | CatalogItem::Func(_)
4049                | CatalogItem::Secret(_)
4050                | CatalogItem::Connection(_) => (),
4051            }
4052        }
4053
4054        Ok(uncached_expressions)
4055    }
4056
4057    /// Selects for each compute dataflow an as-of suitable for bootstrapping it.
4058    ///
4059    /// Returns a set of [`ReadHold`]s that ensures the read frontiers of involved collections stay
4060    /// in place and that must not be dropped before all compute dataflows have been created with
4061    /// the compute controller.
4062    ///
4063    /// This method expects all storage collections and dataflow plans to be available, so it must
4064    /// run after [`Coordinator::bootstrap_storage_collections`] and
4065    /// [`Coordinator::bootstrap_dataflow_plans`].
4066    async fn bootstrap_dataflow_as_ofs(&mut self) -> BTreeMap<GlobalId, ReadHold> {
4067        let mut catalog_ids = Vec::new();
4068        let mut dataflows = Vec::new();
4069        let mut read_policies = BTreeMap::new();
4070        for entry in self.catalog.entries() {
4071            let gid = match entry.item() {
4072                CatalogItem::Index(idx) => idx.global_id(),
4073                CatalogItem::MaterializedView(mv) => mv.global_id_writes(),
4074                CatalogItem::MetricSink(metric_sink) => metric_sink.global_id,
4075                CatalogItem::Table(_)
4076                | CatalogItem::Source(_)
4077                | CatalogItem::Log(_)
4078                | CatalogItem::View(_)
4079                | CatalogItem::Sink(_)
4080                | CatalogItem::Type(_)
4081                | CatalogItem::Func(_)
4082                | CatalogItem::Secret(_)
4083                | CatalogItem::Connection(_) => continue,
4084            };
4085            if let Some(plan) = self.catalog.try_get_physical_plan(&gid) {
4086                catalog_ids.push(gid);
4087                dataflows.push(plan.clone());
4088
4089                if let Some(compaction_window) = entry.item().initial_logical_compaction_window() {
4090                    read_policies.insert(gid, compaction_window.into());
4091                }
4092            }
4093        }
4094
4095        let read_ts = self.get_local_read_ts().await;
4096        let read_holds = as_of_selection::run(
4097            &mut dataflows,
4098            &read_policies,
4099            &*self.controller.storage_collections,
4100            read_ts,
4101            self.controller.read_only(),
4102        );
4103
4104        let catalog = self.catalog_mut();
4105        for (id, plan) in catalog_ids.into_iter().zip_eq(dataflows) {
4106            catalog.set_physical_plan(id, plan);
4107        }
4108
4109        read_holds
4110    }
4111
4112    /// Serves the coordinator, receiving commands from users over `cmd_rx`
4113    /// and feedback from dataflow workers over `feedback_rx`.
4114    ///
4115    /// You must call `bootstrap` before calling this method.
4116    ///
4117    /// BOXED FUTURE: As of Nov 2023 the returned Future from this function was 92KB. This would
4118    /// get stored on the stack which is bad for runtime performance, and blow up our stack usage.
4119    /// Because of that we purposefully move this Future onto the heap (i.e. Box it).
4120    fn serve(
4121        mut self,
4122        mut internal_cmd_rx: mpsc::UnboundedReceiver<Message>,
4123        mut strict_serializable_reads_rx: mpsc::UnboundedReceiver<(ConnectionId, PendingReadTxn)>,
4124        mut cmd_rx: mpsc::UnboundedReceiver<(OpenTelemetryContext, Command)>,
4125        group_commit_rx: appends::GroupCommitWaiter,
4126    ) -> LocalBoxFuture<'static, ()> {
4127        async move {
4128            // Watcher that listens for and reports cluster service status changes.
4129            let mut cluster_events = self.controller.events_stream();
4130            let last_message = Arc::new(Mutex::new(LastMessage {
4131                kind: "none",
4132                stmt: None,
4133            }));
4134
4135            let (idle_tx, mut idle_rx) = tokio::sync::mpsc::channel(1);
4136            let idle_metric = self.metrics.queue_busy_seconds.clone();
4137            let last_message_watchdog = Arc::clone(&last_message);
4138
4139            spawn(|| "coord watchdog", async move {
4140                // Every 5 seconds, attempt to measure how long it takes for the
4141                // coord select loop to be empty, because this message is the last
4142                // processed. If it is idle, this will result in some microseconds
4143                // of measurement.
4144                let mut interval = tokio::time::interval(Duration::from_secs(5));
4145                // If we end up having to wait more than 5 seconds for the coord to respond, then the
4146                // behavior of Delay results in the interval "restarting" from whenever we yield
4147                // instead of trying to catch up.
4148                interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
4149
4150                // Track if we become stuck to de-dupe error reporting.
4151                let mut coord_stuck = false;
4152
4153                loop {
4154                    interval.tick().await;
4155
4156                    // Wait for space in the channel, if we timeout then the coordinator is stuck!
4157                    let duration = tokio::time::Duration::from_secs(30);
4158                    let timeout = tokio::time::timeout(duration, idle_tx.reserve()).await;
4159                    let Ok(maybe_permit) = timeout else {
4160                        // Only log if we're newly stuck, to prevent logging repeatedly.
4161                        if !coord_stuck {
4162                            let last_message = last_message_watchdog.lock().expect("poisoned");
4163                            tracing::warn!(
4164                                last_message_kind = %last_message.kind,
4165                                last_message_sql = %last_message.stmt_to_string(),
4166                                "coordinator stuck for {duration:?}",
4167                            );
4168                        }
4169                        coord_stuck = true;
4170
4171                        continue;
4172                    };
4173
4174                    // We got a permit, we're not stuck!
4175                    if coord_stuck {
4176                        tracing::info!("Coordinator became unstuck");
4177                    }
4178                    coord_stuck = false;
4179
4180                    // If we failed to acquire a permit it's because we're shutting down.
4181                    let Ok(permit) = maybe_permit else {
4182                        break;
4183                    };
4184
4185                    permit.send(idle_metric.start_timer());
4186                }
4187            });
4188
4189            self.schedule_storage_usage_collection().await;
4190            self.schedule_arrangement_sizes_collection().await;
4191            self.schedule_hydration_history_collection();
4192            self.spawn_privatelink_vpc_endpoints_watch_task();
4193            self.spawn_statement_logging_task();
4194            self.spawn_catalog_info_metrics_task();
4195            self.spawn_cluster_controller_task();
4196            flags::tracing_config(self.catalog.system_config()).apply(&self.tracing_handle);
4197
4198            // Report if the handling of a single message takes longer than this threshold.
4199            let warn_threshold = self
4200                .catalog()
4201                .system_config()
4202                .coord_slow_message_warn_threshold();
4203
4204            // How many messages we'd like to batch up before processing them. Must be > 0.
4205            const MESSAGE_BATCH: usize = 64;
4206            let mut messages = Vec::with_capacity(MESSAGE_BATCH);
4207            let mut cmd_messages = Vec::with_capacity(MESSAGE_BATCH);
4208
4209            let message_batch = self.metrics.message_batch.clone();
4210
4211            // A persisted `Notified` future for the linearize re-check signal.
4212            // It must outlive a single loop iteration and be re-`set` only after
4213            // it completes: a fresh `notified()` per iteration could drop a
4214            // wakeup that arrives while a higher-priority branch wins the same
4215            // poll, stranding pending reads. Keeping it pinned across iterations
4216            // leaves it registered, so no wakeup is lost.
4217            let linearize_reads_notify = Arc::clone(&self.linearize_reads_notify);
4218            let linearize_reads_notified = linearize_reads_notify.notified();
4219            tokio::pin!(linearize_reads_notified);
4220
4221            loop {
4222                // Before adding a branch to this select loop, please ensure that the branch is
4223                // cancellation safe and add a comment explaining why. You can refer here for more
4224                // info: https://docs.rs/tokio/latest/tokio/macro.select.html#cancellation-safety
4225                select! {
4226                    // We prioritize internal commands over other commands. However, we work through
4227                    // batches of commands in some branches of this select, which means that even if
4228                    // a command generates internal commands, we will work through the current batch
4229                    // before receiving a new batch of commands.
4230                    biased;
4231
4232                    // `recv_many()` on `UnboundedReceiver` is cancellation safe:
4233                    // https://docs.rs/tokio/1.38.0/tokio/sync/mpsc/struct.UnboundedReceiver.html#cancel-safety-1
4234                    // Receive a batch of commands.
4235                    _ = internal_cmd_rx.recv_many(&mut messages, MESSAGE_BATCH) => {},
4236                    // `next()` on any stream is cancel-safe:
4237                    // https://docs.rs/tokio-stream/0.1.9/tokio_stream/trait.StreamExt.html#cancel-safety
4238                    // Receive a single command.
4239                    Some(event) = cluster_events.next() => {
4240                        messages.push(Message::ClusterEvent(event))
4241                    },
4242                    // See [`mz_controller::Controller::Controller::ready`] for notes
4243                    // on why this is cancel-safe.
4244                    // Receive a single command.
4245                    () = self.controller.ready() => {
4246                        // NOTE: We don't get a `Readiness` back from `ready()`
4247                        // because the controller wants to keep it and it's not
4248                        // trivially `Clone` or `Copy`. Hence this accessor.
4249                        let controller = match self.controller.get_readiness() {
4250                            Readiness::Storage => ControllerReadiness::Storage,
4251                            Readiness::Compute => ControllerReadiness::Compute,
4252                            Readiness::Metrics(_) => ControllerReadiness::Metrics,
4253                            Readiness::Internal(_) => ControllerReadiness::Internal,
4254                            Readiness::NotReady => unreachable!("just signaled as ready"),
4255                        };
4256                        messages.push(Message::ControllerReady { controller });
4257                    }
4258                    // See [`appends::GroupCommitWaiter`] for notes on why this is cancel safe.
4259                    // Receive a single command.
4260                    permit = group_commit_rx.ready() => {
4261                        // If we happen to have batched exactly one user write, use
4262                        // that span so the `emit_trace_id_notice` hooks up.
4263                        // Otherwise, the best we can do is invent a new root span
4264                        // and make it follow from all the Spans in the pending
4265                        // writes.
4266                        let user_write_spans = self.pending_writes.iter().flat_map(|x| match x {
4267                            PendingWriteTxn::User { span, .. } => Some(span),
4268                            PendingWriteTxn::System { .. } => None,
4269                        });
4270                        let span = match user_write_spans.exactly_one() {
4271                            Ok(span) => span.clone(),
4272                            Err(user_write_spans) => {
4273                                let span = info_span!(parent: None, "group_commit_notify");
4274                                for s in user_write_spans {
4275                                    span.follows_from(s);
4276                                }
4277                                span
4278                            }
4279                        };
4280                        messages.push(Message::GroupCommitInitiate(span, Some(permit)));
4281                    },
4282                    // `recv_many()` on `UnboundedReceiver` is cancellation safe:
4283                    // https://docs.rs/tokio/1.38.0/tokio/sync/mpsc/struct.UnboundedReceiver.html#cancel-safety-1
4284                    // Receive a batch of commands.
4285                    count = cmd_rx.recv_many(&mut cmd_messages, MESSAGE_BATCH) => {
4286                        if count == 0 {
4287                            break;
4288                        } else {
4289                            messages.extend(cmd_messages.drain(..).map(
4290                                |(otel_ctx, cmd)| Message::Command(otel_ctx, cmd),
4291                            ));
4292                        }
4293                    },
4294                    // `recv()` on `UnboundedReceiver` is cancellation safe:
4295                    // https://docs.rs/tokio/1.38.0/tokio/sync/mpsc/struct.UnboundedReceiver.html#cancel-safety
4296                    // Receive a single command.
4297                    Some(pending_read_txn) = strict_serializable_reads_rx.recv() => {
4298                        let mut pending_read_txns = vec![pending_read_txn];
4299                        while let Ok(pending_read_txn) = strict_serializable_reads_rx.try_recv() {
4300                            pending_read_txns.push(pending_read_txn);
4301                        }
4302                        for (conn_id, pending_read_txn) in pending_read_txns {
4303                            let prev = self
4304                                .pending_linearize_read_txns
4305                                .insert(conn_id, pending_read_txn);
4306                            soft_assert_or_log!(
4307                                prev.is_none(),
4308                                "connections can not have multiple concurrent reads, prev: {prev:?}"
4309                            )
4310                        }
4311                        messages.push(Message::LinearizeReads);
4312                    }
4313                    // `tick()` on `Interval` is cancel-safe:
4314                    // https://docs.rs/tokio/1.19.2/tokio/time/struct.Interval.html#cancel-safety
4315                    // Receive a single command.
4316                    _ = self.advance_timelines_interval.tick() => {
4317                        // Writable keepalives use the committer to advance tables and read holds.
4318                        // Its permit coalesces ticks behind a slow oracle. Read-only mode advances
4319                        // timelines directly.
4320                        if self.controller.read_only() {
4321                            messages.push(Message::AdvanceTimelines);
4322                        } else {
4323                            self.group_commit_tx.notify();
4324                        }
4325                    },
4326                    // Re-check pending strict serializable reads. Deliberately
4327                    // placed below the group commit branches above: a re-check
4328                    // only makes a read ready if the timestamp oracle has
4329                    // advanced, and the oracle only advances via group commit, so
4330                    // this must never win over (and thereby starve) group commit.
4331                    // `Notify` coalesces re-arms into a single wakeup, so even
4332                    // when a pending read sits just behind the oracle (re-armed
4333                    // sub-millisecond), the lower branches (including the idle
4334                    // watchdog) stay reachable. See the pin above for why the
4335                    // future is persisted rather than recreated per iteration.
4336                    () = linearize_reads_notified.as_mut() => {
4337                        linearize_reads_notified.set(linearize_reads_notify.notified());
4338                        messages.push(Message::LinearizeReads);
4339                    }
4340                    // `tick()` on `Interval` is cancel-safe:
4341                    // https://docs.rs/tokio/1.19.2/tokio/time/struct.Interval.html#cancel-safety
4342                    // Receive a single command.
4343                    _ = self.caught_up_check_interval.tick() => {
4344                        // We do this directly on the main loop instead of
4345                        // firing off a message. We are still in read-only mode,
4346                        // so optimizing for latency, not blocking the main loop
4347                        // is not that important.
4348                        self.maybe_check_caught_up().await;
4349
4350                        continue;
4351                    },
4352
4353                    // Process the idle metric at the lowest priority to sample queue non-idle time.
4354                    // `recv()` on `Receiver` is cancellation safe:
4355                    // https://docs.rs/tokio/1.8.0/tokio/sync/mpsc/struct.Receiver.html#cancel-safety
4356                    // Receive a single command.
4357                    timer = idle_rx.recv() => {
4358                        timer.expect("does not drop").observe_duration();
4359                        self.metrics
4360                            .message_handling
4361                            .with_label_values(&["watchdog"])
4362                            .observe(0.0);
4363                        continue;
4364                    }
4365                };
4366
4367                // Observe the number of messages we're processing at once.
4368                message_batch.observe(f64::cast_lossy(messages.len()));
4369
4370                for msg in messages.drain(..) {
4371                    // All message processing functions trace. Start a parent span
4372                    // for them to make it easy to find slow messages.
4373                    let msg_kind = msg.kind();
4374                    let span = span!(
4375                        target: "mz_adapter::coord::handle_message_loop",
4376                        Level::INFO,
4377                        "coord::handle_message",
4378                        kind = msg_kind
4379                    );
4380                    let otel_context = span.context().span().span_context().clone();
4381
4382                    // Record the last kind of message in case we get stuck. For
4383                    // execute commands, we additionally stash the user's SQL,
4384                    // statement, so we can log it in case we get stuck.
4385                    *last_message.lock().expect("poisoned") = LastMessage {
4386                        kind: msg_kind,
4387                        stmt: match &msg {
4388                            Message::Command(
4389                                _,
4390                                Command::Execute {
4391                                    portal_name,
4392                                    session,
4393                                    ..
4394                                },
4395                            ) => session
4396                                .get_portal_unverified(portal_name)
4397                                .and_then(|p| p.stmt.as_ref().map(Arc::clone)),
4398                            _ => None,
4399                        },
4400                    };
4401
4402                    let start = Instant::now();
4403                    self.handle_message(msg).instrument(span).await;
4404                    let duration = start.elapsed();
4405
4406                    self.metrics
4407                        .message_handling
4408                        .with_label_values(&[msg_kind])
4409                        .observe(duration.as_secs_f64());
4410
4411                    // If something is _really_ slow, print a trace id for debugging, if OTEL is enabled.
4412                    if duration > warn_threshold {
4413                        let trace_id = otel_context.is_valid().then(|| otel_context.trace_id());
4414                        tracing::error!(
4415                            ?msg_kind,
4416                            ?trace_id,
4417                            ?duration,
4418                            "very slow coordinator message"
4419                        );
4420                    }
4421                }
4422            }
4423
4424            // The sweep can own timestamp-oracle senders through its background
4425            // client. Release them before the coordinator runtime starts shutting
4426            // down the oracle workers.
4427            if let Some(sweep) = self.hydration_history_sweep.take() {
4428                sweep.abort_and_wait().await;
4429            }
4430
4431            // Try and cleanup as a best effort. There may be some async tasks out there holding a
4432            // reference that prevents us from cleaning up.
4433            if let Some(catalog) = Arc::into_inner(self.catalog) {
4434                catalog.expire().await;
4435            }
4436        }
4437        .boxed_local()
4438    }
4439
4440    /// Obtain a read-only Catalog reference.
4441    fn catalog(&self) -> &Catalog {
4442        &self.catalog
4443    }
4444
4445    /// Obtain a read-only Catalog snapshot, suitable for giving out to
4446    /// non-Coordinator thread tasks.
4447    fn owned_catalog(&self) -> Arc<Catalog> {
4448        Arc::clone(&self.catalog)
4449    }
4450
4451    /// Obtain a handle to the optimizer metrics, suitable for giving
4452    /// out to non-Coordinator thread tasks.
4453    fn optimizer_metrics(&self) -> OptimizerMetrics {
4454        self.optimizer_metrics.clone()
4455    }
4456
4457    /// Obtain a writeable Catalog reference.
4458    fn catalog_mut(&mut self) -> &mut Catalog {
4459        // make_mut will cause any other Arc references (from owned_catalog) to
4460        // continue to be valid by cloning the catalog, putting it in a new Arc,
4461        // which lives at self._catalog. If there are no other Arc references,
4462        // then no clone is made, and it returns a reference to the existing
4463        // object. This makes this method and owned_catalog both very cheap: at
4464        // most one clone per catalog mutation, but only if there's a read-only
4465        // reference to it.
4466        Arc::make_mut(&mut self.catalog)
4467    }
4468
4469    /// Refills the user ID pool by allocating IDs from the catalog.
4470    ///
4471    /// Requests `max(min_count, batch_size)` IDs so the pool is never
4472    /// under-filled relative to the configured batch size.
4473    async fn refill_user_id_pool(&mut self, min_count: u64) -> Result<(), AdapterError> {
4474        let batch_size = USER_ID_POOL_BATCH_SIZE.get(self.catalog().system_config().dyncfgs());
4475        let to_allocate = min_count.max(u64::from(batch_size));
4476        let id_ts = self.get_catalog_write_ts().await;
4477        let ids = self.catalog().allocate_user_ids(to_allocate, id_ts).await?;
4478        if let (Some((first_id, _)), Some((last_id, _))) = (ids.first(), ids.last()) {
4479            let start = match first_id {
4480                CatalogItemId::User(id) => *id,
4481                other => {
4482                    return Err(AdapterError::Internal(format!(
4483                        "expected User CatalogItemId, got {other:?}"
4484                    )));
4485                }
4486            };
4487            let end = match last_id {
4488                CatalogItemId::User(id) => *id + 1, // exclusive upper bound
4489                other => {
4490                    return Err(AdapterError::Internal(format!(
4491                        "expected User CatalogItemId, got {other:?}"
4492                    )));
4493                }
4494            };
4495            self.user_id_pool.refill(start, end);
4496        } else {
4497            return Err(AdapterError::Internal(
4498                "catalog returned no user IDs".into(),
4499            ));
4500        }
4501        Ok(())
4502    }
4503
4504    /// Allocates a single user ID, refilling the pool from the catalog if needed.
4505    async fn allocate_user_id(&mut self) -> Result<(CatalogItemId, GlobalId), AdapterError> {
4506        if let Some(id) = self.user_id_pool.allocate() {
4507            return Ok((CatalogItemId::User(id), GlobalId::User(id)));
4508        }
4509        self.refill_user_id_pool(1).await?;
4510        let id = self.user_id_pool.allocate().expect("ID pool just refilled");
4511        Ok((CatalogItemId::User(id), GlobalId::User(id)))
4512    }
4513
4514    /// Allocates `count` user IDs, refilling the pool from the catalog if needed.
4515    async fn allocate_user_ids(
4516        &mut self,
4517        count: u64,
4518    ) -> Result<Vec<(CatalogItemId, GlobalId)>, AdapterError> {
4519        if self.user_id_pool.remaining() < count {
4520            self.refill_user_id_pool(count).await?;
4521        }
4522        let raw_ids = self
4523            .user_id_pool
4524            .allocate_many(count)
4525            .expect("pool has enough IDs after refill");
4526        Ok(raw_ids
4527            .into_iter()
4528            .map(|id| (CatalogItemId::User(id), GlobalId::User(id)))
4529            .collect())
4530    }
4531
4532    /// Obtain a reference to the coordinator's connection context.
4533    fn connection_context(&self) -> &ConnectionContext {
4534        self.controller.connection_context()
4535    }
4536
4537    /// Obtain a reference to the coordinator's secret reader, in an `Arc`.
4538    fn secrets_reader(&self) -> &Arc<dyn SecretsReader> {
4539        &self.connection_context().secrets_reader
4540    }
4541
4542    /// Publishes a notice message to all sessions.
4543    ///
4544    /// TODO(parkmycar): This code is dead, but is a nice parallel to [`Coordinator::broadcast_notice_tx`]
4545    /// so we keep it around.
4546    #[allow(dead_code)]
4547    pub(crate) fn broadcast_notice(&self, notice: AdapterNotice) {
4548        for meta in self.active_conns.values() {
4549            let _ = meta.notice_tx.send(notice.clone());
4550        }
4551    }
4552
4553    /// Returns a closure that will publish a notice to all sessions that were active at the time
4554    /// this method was called.
4555    pub(crate) fn broadcast_notice_tx(
4556        &self,
4557    ) -> Box<dyn FnOnce(AdapterNotice) -> () + Send + 'static> {
4558        let senders: Vec<_> = self
4559            .active_conns
4560            .values()
4561            .map(|meta| meta.notice_tx.clone())
4562            .collect();
4563        Box::new(move |notice| {
4564            for tx in senders {
4565                let _ = tx.send(notice.clone());
4566            }
4567        })
4568    }
4569
4570    pub(crate) fn active_conns(&self) -> &BTreeMap<ConnectionId, ConnMeta> {
4571        &self.active_conns
4572    }
4573
4574    #[instrument(level = "debug")]
4575    pub(crate) fn retire_execution(
4576        &mut self,
4577        reason: StatementEndedExecutionReason,
4578        ctx_extra: ExecuteContextExtra,
4579    ) {
4580        if let Some(uuid) = ctx_extra.retire() {
4581            let ended_at = self.now();
4582            self.end_statement_execution(uuid, reason, ended_at);
4583        }
4584    }
4585
4586    /// Creates a new dataflow builder from the catalog and indexes in `self`.
4587    #[instrument(level = "debug")]
4588    pub fn dataflow_builder(&self, instance: ComputeInstanceId) -> DataflowBuilder<'_> {
4589        let compute = self
4590            .instance_snapshot(instance)
4591            .expect("compute instance does not exist");
4592        DataflowBuilder::new(self.catalog().state(), compute)
4593    }
4594
4595    /// Return a reference-less snapshot to the indicated compute instance.
4596    pub fn instance_snapshot(
4597        &self,
4598        id: ComputeInstanceId,
4599    ) -> Result<ComputeInstanceSnapshot, InstanceMissing> {
4600        ComputeInstanceSnapshot::new(&self.controller, id)
4601    }
4602
4603    /// Call into the compute controller to install a finalized dataflow, and
4604    /// initialize the read policies for its exported readable objects.
4605    ///
4606    /// # Panics
4607    ///
4608    /// Panics if dataflow creation fails.
4609    pub(crate) async fn ship_dataflow(
4610        &mut self,
4611        dataflow: DataflowDescription<LirRelationExpr>,
4612        instance: ComputeInstanceId,
4613        target_replica: Option<ReplicaId>,
4614    ) {
4615        self.try_ship_dataflow(dataflow, instance, target_replica)
4616            .await
4617            .unwrap_or_terminate("dataflow creation cannot fail");
4618    }
4619
4620    /// Call into the compute controller to install a finalized dataflow, and
4621    /// initialize the read policies for its exported readable objects.
4622    pub(crate) async fn try_ship_dataflow(
4623        &mut self,
4624        dataflow: DataflowDescription<LirRelationExpr>,
4625        instance: ComputeInstanceId,
4626        target_replica: Option<ReplicaId>,
4627    ) -> Result<(), DataflowCreationError> {
4628        // We must only install read policies for indexes, not for sinks.
4629        // Sinks are write-only compute collections that don't have read policies.
4630        let export_ids = dataflow.exported_index_ids().collect();
4631
4632        self.controller
4633            .compute
4634            .create_dataflow(instance, dataflow, target_replica)?;
4635
4636        self.initialize_compute_read_policies(export_ids, instance, CompactionWindow::Default)
4637            .await;
4638
4639        Ok(())
4640    }
4641
4642    /// Call into the compute controller to allow writes to the specified IDs
4643    /// from the specified instance. Calling this function multiple times and
4644    /// calling it on a read-only instance has no effect.
4645    pub(crate) fn allow_writes(&mut self, instance: ComputeInstanceId, id: GlobalId) {
4646        self.controller
4647            .compute
4648            .allow_writes(instance, id)
4649            .unwrap_or_terminate("allow_writes cannot fail");
4650    }
4651
4652    /// Like `ship_dataflow`, but also await on builtin table updates.
4653    pub(crate) async fn ship_dataflow_and_notice_builtin_table_updates(
4654        &mut self,
4655        dataflow: DataflowDescription<LirRelationExpr>,
4656        instance: ComputeInstanceId,
4657        notice_builtin_updates_fut: Option<BuiltinTableAppendNotify>,
4658        target_replica: Option<ReplicaId>,
4659    ) {
4660        if let Some(notice_builtin_updates_fut) = notice_builtin_updates_fut {
4661            let ship_dataflow_fut = self.ship_dataflow(dataflow, instance, target_replica);
4662            let ((), ()) =
4663                futures::future::join(notice_builtin_updates_fut, ship_dataflow_fut).await;
4664        } else {
4665            self.ship_dataflow(dataflow, instance, target_replica).await;
4666        }
4667    }
4668
4669    /// Install a _watch set_ in the controller that is automatically associated with the given
4670    /// connection id. The watchset will be automatically cleared if the connection terminates
4671    /// before the watchset completes.
4672    pub fn install_compute_watch_set(
4673        &mut self,
4674        conn_id: ConnectionId,
4675        objects: BTreeSet<GlobalId>,
4676        t: Timestamp,
4677        state: WatchSetResponse,
4678    ) -> Result<(), CollectionLookupError> {
4679        let ws_id = self.controller.install_compute_watch_set(objects, t)?;
4680        self.connection_watch_sets
4681            .entry(conn_id.clone())
4682            .or_default()
4683            .insert(ws_id);
4684        self.installed_watch_sets.insert(ws_id, (conn_id, state));
4685        Ok(())
4686    }
4687
4688    /// Install a _watch set_ in the controller that is automatically associated with the given
4689    /// connection id. The watchset will be automatically cleared if the connection terminates
4690    /// before the watchset completes.
4691    pub fn install_storage_watch_set(
4692        &mut self,
4693        conn_id: ConnectionId,
4694        objects: BTreeSet<GlobalId>,
4695        t: Timestamp,
4696        state: WatchSetResponse,
4697    ) -> Result<(), CollectionMissing> {
4698        let ws_id = self.controller.install_storage_watch_set(objects, t)?;
4699        self.connection_watch_sets
4700            .entry(conn_id.clone())
4701            .or_default()
4702            .insert(ws_id);
4703        self.installed_watch_sets.insert(ws_id, (conn_id, state));
4704        Ok(())
4705    }
4706
4707    /// Cancels pending watchsets associated with the provided connection id.
4708    pub fn cancel_pending_watchsets(&mut self, conn_id: &ConnectionId) {
4709        if let Some(ws_ids) = self.connection_watch_sets.remove(conn_id) {
4710            for ws_id in ws_ids {
4711                self.installed_watch_sets.remove(&ws_id);
4712            }
4713        }
4714    }
4715
4716    /// Returns the state of the [`Coordinator`] formatted as JSON.
4717    ///
4718    /// The returned value is not guaranteed to be stable and may change at any point in time.
4719    pub async fn dump(&self) -> Result<serde_json::Value, anyhow::Error> {
4720        // Note: We purposefully use the `Debug` formatting for the value of all fields in the
4721        // returned object as a tradeoff between usability and stability. `serde_json` will fail
4722        // to serialize an object if the keys aren't strings, so `Debug` formatting the values
4723        // prevents a future unrelated change from silently breaking this method.
4724
4725        let global_timelines: BTreeMap<_, _> = self
4726            .global_timelines
4727            .iter()
4728            .map(|(timeline, state)| (timeline.to_string(), format!("{state:?}")))
4729            .collect();
4730        let active_conns: BTreeMap<_, _> = self
4731            .active_conns
4732            .iter()
4733            .map(|(id, meta)| (id.unhandled().to_string(), format!("{meta:?}")))
4734            .collect();
4735        let txn_read_holds: BTreeMap<_, _> = self
4736            .txn_read_holds
4737            .iter()
4738            .map(|(id, capability)| (id.unhandled().to_string(), format!("{capability:?}")))
4739            .collect();
4740        let pending_peeks: BTreeMap<_, _> = self
4741            .pending_peeks
4742            .iter()
4743            .map(|(id, peek)| (id.to_string(), format!("{peek:?}")))
4744            .collect();
4745        let client_pending_peeks: BTreeMap<_, _> = self
4746            .client_pending_peeks
4747            .iter()
4748            .map(|(id, peek)| {
4749                let peek: BTreeMap<_, _> = peek
4750                    .iter()
4751                    .map(|(uuid, storage_id)| (uuid.to_string(), storage_id))
4752                    .collect();
4753                (id.to_string(), peek)
4754            })
4755            .collect();
4756        let pending_linearize_read_txns: BTreeMap<_, _> = self
4757            .pending_linearize_read_txns
4758            .iter()
4759            .map(|(id, read_txn)| (id.unhandled().to_string(), format!("{read_txn:?}")))
4760            .collect();
4761
4762        Ok(serde_json::json!({
4763            "global_timelines": global_timelines,
4764            "active_conns": active_conns,
4765            "txn_read_holds": txn_read_holds,
4766            "pending_peeks": pending_peeks,
4767            "client_pending_peeks": client_pending_peeks,
4768            "pending_linearize_read_txns": pending_linearize_read_txns,
4769            "controller": self.controller.dump().await?,
4770        }))
4771    }
4772
4773    /// Prune all storage usage events from the [`MZ_STORAGE_USAGE_BY_SHARD`] table that are older
4774    /// than `retention_period`.
4775    ///
4776    /// This method will read the entire contents of [`MZ_STORAGE_USAGE_BY_SHARD`] into memory
4777    /// which can be expensive.
4778    ///
4779    /// DO NOT call this method outside of startup. The safety of reading at the current oracle read
4780    /// timestamp and then writing at whatever the current write timestamp is (instead of
4781    /// `read_ts + 1`) relies on the fact that there are no outstanding writes during startup.
4782    ///
4783    /// Group commit, which this method uses to write the retractions, has builtin fencing, and we
4784    /// never commit retractions to [`MZ_STORAGE_USAGE_BY_SHARD`] outside of this method, which is
4785    /// only called once during startup. So we don't have to worry about double/invalid retractions.
4786    async fn prune_storage_usage_events_on_startup(&self, retention_period: Duration) {
4787        let item_id = self
4788            .catalog()
4789            .resolve_builtin_table(&MZ_STORAGE_USAGE_BY_SHARD);
4790        let global_id = self.catalog.get_entry(&item_id).latest_global_id();
4791        let read_ts = self.get_local_read_ts().await;
4792        let current_contents_fut = self
4793            .controller
4794            .storage_collections
4795            .snapshot(global_id, read_ts);
4796        let internal_cmd_tx = self.internal_cmd_tx.clone();
4797        spawn(|| "storage_usage_prune", async move {
4798            let mut current_contents = current_contents_fut
4799                .await
4800                .unwrap_or_terminate("cannot fail to fetch snapshot");
4801            differential_dataflow::consolidation::consolidate(&mut current_contents);
4802
4803            let cutoff_ts = u128::from(read_ts).saturating_sub(retention_period.as_millis());
4804            let mut expired = Vec::new();
4805            for (row, diff) in current_contents {
4806                assert_eq!(
4807                    diff, 1,
4808                    "consolidated contents should not contain retractions: ({row:#?}, {diff:#?})"
4809                );
4810                // This logic relies on the definition of `mz_storage_usage_by_shard` not changing.
4811                let collection_timestamp = row
4812                    .unpack()
4813                    .get(3)
4814                    .expect("definition of mz_storage_by_shard changed")
4815                    .unwrap_timestamptz();
4816                let collection_timestamp = collection_timestamp.timestamp_millis();
4817                let collection_timestamp: u128 = collection_timestamp
4818                    .try_into()
4819                    .expect("all collections happen after Jan 1 1970");
4820                if collection_timestamp < cutoff_ts {
4821                    debug!("pruning storage event {row:?}");
4822                    let builtin_update = BuiltinTableUpdate::row(item_id, row, Diff::MINUS_ONE);
4823                    expired.push(builtin_update);
4824                }
4825            }
4826
4827            // main thread has shut down.
4828            let _ = internal_cmd_tx.send(Message::StorageUsagePrune(expired));
4829        });
4830    }
4831
4832    /// Retracts `mz_object_arrangement_size_history` rows older than the
4833    /// `arrangement_size_history_retention_period` dyncfg.
4834    ///
4835    /// Must only run at startup: it reads at the oracle read timestamp and
4836    /// writes retractions at the current write timestamp, which is only safe
4837    /// when no other writes are in flight. See [the equivalent storage-usage
4838    /// pruner](Self::prune_storage_usage_events_on_startup) for the same
4839    /// reasoning.
4840    async fn prune_arrangement_sizes_history_on_startup(&self) {
4841        // The catalog server is not writable in read-only mode.
4842        if self.controller.read_only() {
4843            return;
4844        }
4845
4846        let retention_period = mz_adapter_types::dyncfgs::ARRANGEMENT_SIZE_HISTORY_RETENTION_PERIOD
4847            .get(self.catalog().system_config().dyncfgs());
4848        let item_id = self
4849            .catalog()
4850            .resolve_builtin_table(&mz_catalog::builtin::MZ_OBJECT_ARRANGEMENT_SIZE_HISTORY);
4851        let global_id = self.catalog.get_entry(&item_id).latest_global_id();
4852        let read_ts = self.get_local_read_ts().await;
4853        let current_contents_fut = self
4854            .controller
4855            .storage_collections
4856            .snapshot(global_id, read_ts);
4857        let internal_cmd_tx = self.internal_cmd_tx.clone();
4858        spawn(|| "arrangement_sizes_history_prune", async move {
4859            let mut current_contents = current_contents_fut
4860                .await
4861                .unwrap_or_terminate("cannot fail to fetch snapshot");
4862            differential_dataflow::consolidation::consolidate(&mut current_contents);
4863
4864            let cutoff_ts = u128::from(read_ts).saturating_sub(retention_period.as_millis());
4865            let expired =
4866                arrangement_sizes_expired_retractions(current_contents, cutoff_ts, item_id);
4867
4868            // TODO(arrangement-sizes): when the writeable-catalog-server
4869            // plumbing in https://github.com/MaterializeInc/materialize/pull/35436
4870            // lands, retract directly on `mz_catalog_server`.
4871            let _ = internal_cmd_tx.send(Message::ArrangementSizesPrune(expired));
4872        });
4873    }
4874
4875    /// The environment's current credit consumption rate, summed over all user
4876    /// cluster replicas except those of `exclude_cluster`.
4877    fn current_credit_consumption_rate(&self, exclude_cluster: Option<ClusterId>) -> Numeric {
4878        self.catalog()
4879            .user_cluster_replicas()
4880            .filter(|replica| Some(replica.cluster_id) != exclude_cluster)
4881            .filter_map(|replica| match &replica.config.location {
4882                ReplicaLocation::Managed(location) => Some(self.replica_credits_per_hour(location)),
4883                ReplicaLocation::Unmanaged(_) => None,
4884            })
4885            .sum()
4886    }
4887
4888    /// The credit rate of a managed replica, read from the size map by its billing size.
4889    ///
4890    /// An unknown billing size counts as free. DDL validates `SIZE` and `BILLED AS` against
4891    /// the map at replica creation, but the map is external configuration and can lose a
4892    /// size later. That case is a soft panic rather than a hard one, so that in production
4893    /// such a replica can still be dropped from SQL.
4894    fn replica_credits_per_hour(&self, location: &ManagedReplicaLocation) -> Numeric {
4895        let size = location.size_for_billing();
4896        match self.catalog().cluster_replica_sizes().0.get(size) {
4897            Some(allocation) => allocation.credits_per_hour,
4898            None => {
4899                soft_panic_or_log!(
4900                    "replica of size {:?} bills as unknown replica size {:?}, counting it as free",
4901                    location.size,
4902                    size,
4903                );
4904                Numeric::zero()
4905            }
4906        }
4907    }
4908}
4909
4910/// Returns retraction updates for rows in a consolidated
4911/// `mz_object_arrangement_size_history` snapshot whose `collection_timestamp`
4912/// (column 3) is strictly before `cutoff_ts`.
4913///
4914/// Panics if any input row has `diff != 1`: the caller must consolidate first,
4915/// and a consolidated history table should never contain retractions because
4916/// the only source of retractions is this function itself.
4917fn arrangement_sizes_expired_retractions(
4918    rows: impl IntoIterator<Item = (mz_repr::Row, i64)>,
4919    cutoff_ts: u128,
4920    item_id: CatalogItemId,
4921) -> Vec<BuiltinTableUpdate> {
4922    let mut expired = Vec::new();
4923    for (row, diff) in rows {
4924        assert_eq!(
4925            diff, 1,
4926            "consolidated contents should not contain retractions: ({row:#?}, {diff:#?})"
4927        );
4928        let collection_timestamp = row
4929            .unpack()
4930            .get(3)
4931            .expect("definition of mz_object_arrangement_size_history changed")
4932            .unwrap_timestamptz()
4933            .timestamp_millis();
4934        let collection_timestamp: u128 = collection_timestamp
4935            .try_into()
4936            .expect("all collections happen after Jan 1 1970");
4937        if collection_timestamp < cutoff_ts {
4938            expired.push(BuiltinTableUpdate::row(item_id, row, Diff::MINUS_ONE));
4939        }
4940    }
4941    expired
4942}
4943
4944#[cfg(test)]
4945impl Coordinator {
4946    #[allow(dead_code)]
4947    async fn verify_ship_dataflow_no_error(
4948        &mut self,
4949        dataflow: DataflowDescription<LirRelationExpr>,
4950    ) {
4951        // `ship_dataflow_new` is not allowed to have a `Result` return because this function is
4952        // called after `catalog_transact`, after which no errors are allowed. This test exists to
4953        // prevent us from incorrectly teaching those functions how to return errors (which has
4954        // happened twice and is the motivation for this test).
4955
4956        // An arbitrary compute instance ID to satisfy the function calls below. Note that
4957        // this only works because this function will never run.
4958        let compute_instance = ComputeInstanceId::user(1).expect("1 is a valid ID");
4959
4960        let _: () = self.ship_dataflow(dataflow, compute_instance, None).await;
4961    }
4962}
4963
4964/// Contains information about the last message the [`Coordinator`] processed.
4965struct LastMessage {
4966    kind: &'static str,
4967    stmt: Option<Arc<Statement<Raw>>>,
4968}
4969
4970impl LastMessage {
4971    /// Returns a redacted version of the statement that is safe for logs.
4972    fn stmt_to_string(&self) -> Cow<'static, str> {
4973        self.stmt
4974            .as_ref()
4975            .map(|stmt| stmt.to_ast_string_redacted().into())
4976            .unwrap_or(Cow::Borrowed("<none>"))
4977    }
4978}
4979
4980impl fmt::Debug for LastMessage {
4981    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
4982        f.debug_struct("LastMessage")
4983            .field("kind", &self.kind)
4984            .field("stmt", &self.stmt_to_string())
4985            .finish()
4986    }
4987}
4988
4989impl Drop for LastMessage {
4990    fn drop(&mut self) {
4991        // Only print the last message if we're currently panicking, otherwise we'd spam our logs.
4992        if std::thread::panicking() {
4993            // If we're panicking theres no guarantee `tracing` still works, so print to stderr.
4994            eprintln!("Coordinator panicking, dumping last message\n{self:?}",);
4995        }
4996    }
4997}
4998
4999/// Serves the coordinator based on the provided configuration.
5000///
5001/// For a high-level description of the coordinator, see the [crate
5002/// documentation](crate).
5003///
5004/// Returns a handle to the coordinator and a client to communicate with the
5005/// coordinator.
5006///
5007/// BOXED FUTURE: As of Nov 2023 the returned Future from this function was 42KB. This would
5008/// get stored on the stack which is bad for runtime performance, and blow up our stack usage.
5009/// Because of that we purposefully move this Future onto the heap (i.e. Box it).
5010pub fn serve(
5011    Config {
5012        controller_config,
5013        controller_envd_epoch,
5014        mut storage,
5015        timestamp_oracle_url,
5016        unsafe_mode,
5017        all_features,
5018        build_info,
5019        environment_id,
5020        metrics_registry,
5021        now,
5022        secrets_controller,
5023        cloud_resource_controller,
5024        cluster_replica_sizes,
5025        builtin_system_cluster_config,
5026        builtin_catalog_server_cluster_config,
5027        builtin_probe_cluster_config,
5028        builtin_support_cluster_config,
5029        builtin_analytics_cluster_config,
5030        system_parameter_defaults,
5031        availability_zones,
5032        storage_usage_client,
5033        storage_usage_collection_interval,
5034        storage_usage_retention_period,
5035        segment_client,
5036        egress_addresses,
5037        aws_account_id,
5038        aws_privatelink_availability_zones,
5039        connection_context,
5040        connection_limit_callback,
5041        remote_system_parameters,
5042        webhook_concurrency_limit,
5043        http_host_name,
5044        tracing_handle,
5045        read_only_controllers,
5046        caught_up_trigger: clusters_caught_up_trigger,
5047        helm_chart_version,
5048        license_key,
5049        external_login_password_mz_system,
5050        force_builtin_schema_migration,
5051    }: Config,
5052) -> BoxFuture<'static, Result<(Handle, Client), AdapterError>> {
5053    async move {
5054        let coord_start = Instant::now();
5055        info!("startup: coordinator init: beginning");
5056        info!("startup: coordinator init: preamble beginning");
5057
5058        // Initializing the builtins can be an expensive process and consume a lot of memory. We
5059        // forcibly initialize it early while the stack is relatively empty to avoid stack
5060        // overflows later.
5061        let _builtins = LazyLock::force(&BUILTINS_STATIC);
5062
5063        let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
5064        let (internal_cmd_tx, internal_cmd_rx) = mpsc::unbounded_channel();
5065        let (strict_serializable_reads_tx, strict_serializable_reads_rx) =
5066            mpsc::unbounded_channel();
5067
5068        // Validate and process availability zones.
5069        if !availability_zones.iter().all_unique() {
5070            coord_bail!("availability zones must be unique");
5071        }
5072
5073        let aws_principal_context = match (
5074            aws_account_id,
5075            connection_context.aws_external_id_prefix.clone(),
5076        ) {
5077            (Some(aws_account_id), Some(aws_external_id_prefix)) => Some(AwsPrincipalContext {
5078                aws_account_id,
5079                aws_external_id_prefix,
5080            }),
5081            _ => None,
5082        };
5083
5084        let aws_privatelink_availability_zones = aws_privatelink_availability_zones
5085            .map(|azs_vec| BTreeSet::from_iter(azs_vec.iter().cloned()));
5086
5087        info!(
5088            "startup: coordinator init: preamble complete in {:?}",
5089            coord_start.elapsed()
5090        );
5091        let oracle_init_start = Instant::now();
5092        info!("startup: coordinator init: timestamp oracle init beginning");
5093
5094        let timestamp_oracle_config = timestamp_oracle_url
5095            .map(|url| TimestampOracleConfig::from_url(&url, &metrics_registry))
5096            .transpose()?;
5097        let mut initial_timestamps =
5098            get_initial_oracle_timestamps(&timestamp_oracle_config).await?;
5099
5100        // Insert an entry for the `EpochMilliseconds` timeline if one doesn't exist,
5101        // which will ensure that the timeline is initialized since it's required
5102        // by the system.
5103        initial_timestamps
5104            .entry(Timeline::EpochMilliseconds)
5105            .or_insert_with(mz_repr::Timestamp::minimum);
5106        let mut timestamp_oracles = BTreeMap::new();
5107        for (timeline, initial_timestamp) in initial_timestamps {
5108            Coordinator::ensure_timeline_state_with_initial_time(
5109                &timeline,
5110                initial_timestamp,
5111                now.clone(),
5112                timestamp_oracle_config.clone(),
5113                &mut timestamp_oracles,
5114                read_only_controllers,
5115            )
5116            .await;
5117        }
5118
5119        // Opening the durable catalog uses one or more timestamps without communicating with
5120        // the timestamp oracle. Here we make sure to apply the catalog upper with the timestamp
5121        // oracle to linearize future operations with opening the catalog.
5122        let catalog_upper = storage.current_upper().await;
5123        // Choose a time at which to boot. This is used, for example, to prune
5124        // old storage usage data or migrate audit log entries.
5125        //
5126        // This time is usually the current system time, but with protection
5127        // against backwards time jumps, even across restarts.
5128        let epoch_millis_oracle = &timestamp_oracles
5129            .get(&Timeline::EpochMilliseconds)
5130            .expect("inserted above")
5131            .oracle;
5132
5133        // The catalog shard's upper is durable, so a write that once landed far ahead of the
5134        // clock is re-applied to the oracle here on every boot and cannot be waited out. We
5135        // report it rather than refusing to start: the timeline is stalled either way, and a
5136        // process that will not boot turns that into a total outage plus a crash loop.
5137        let boot_now: mz_repr::Timestamp = (now)().into();
5138        if catalog_upper > timeline::write_ts_upper_bound(&boot_now) {
5139            tracing::error!(
5140                %catalog_upper, %boot_now,
5141                "catalog upper is far ahead of the wall clock, so writes and \
5142                strict-serializable reads on the EpochMilliseconds timeline will block \
5143                until the clock catches up",
5144            );
5145        }
5146
5147        let mut boot_ts = if read_only_controllers {
5148            let read_ts = epoch_millis_oracle.read_ts().await;
5149            std::cmp::max(read_ts, catalog_upper)
5150        } else {
5151            // Getting/applying a write timestamp bumps the write timestamp in the
5152            // oracle, which we're not allowed in read-only mode.
5153            epoch_millis_oracle.apply_write(catalog_upper).await;
5154            epoch_millis_oracle.write_ts().await.timestamp
5155        };
5156
5157        info!(
5158            "startup: coordinator init: timestamp oracle init complete in {:?}",
5159            oracle_init_start.elapsed()
5160        );
5161
5162        let catalog_open_start = Instant::now();
5163        info!("startup: coordinator init: catalog open beginning");
5164        let persist_client = controller_config
5165            .persist_clients
5166            .open(controller_config.persist_location.clone())
5167            .await
5168            .context("opening persist client")?;
5169        let builtin_item_migration_config =
5170            BuiltinItemMigrationConfig {
5171                persist_client: persist_client.clone(),
5172                read_only: read_only_controllers,
5173                force_migration: force_builtin_schema_migration,
5174            }
5175        ;
5176        let OpenCatalogResult {
5177            mut catalog,
5178            last_seen_version,
5179            migrated_storage_collections_0dt,
5180            new_builtin_collections,
5181            builtin_table_updates,
5182            cached_global_exprs,
5183            uncached_local_exprs,
5184        } = Catalog::open(mz_catalog::config::Config {
5185            storage,
5186            metrics_registry: &metrics_registry,
5187            state: mz_catalog::config::StateConfig {
5188                unsafe_mode,
5189                all_features,
5190                build_info,
5191                environment_id: environment_id.clone(),
5192                read_only: read_only_controllers,
5193                now: now.clone(),
5194                boot_ts: boot_ts.clone(),
5195                skip_migrations: false,
5196                cluster_replica_sizes,
5197                builtin_system_cluster_config,
5198                builtin_catalog_server_cluster_config,
5199                builtin_probe_cluster_config,
5200                builtin_support_cluster_config,
5201                builtin_analytics_cluster_config,
5202                system_parameter_defaults,
5203                remote_system_parameters,
5204                availability_zones,
5205                egress_addresses,
5206                aws_principal_context,
5207                aws_privatelink_availability_zones,
5208                connection_context,
5209                http_host_name,
5210                builtin_item_migration_config,
5211                persist_client: persist_client.clone(),
5212                enable_expression_cache_override: None,
5213                helm_chart_version,
5214                external_login_password_mz_system,
5215                license_key: license_key.clone(),
5216            },
5217        })
5218        .await?;
5219
5220        // Opening the catalog uses one or more timestamps, so push the boot timestamp up to the
5221        // current catalog upper.
5222        let catalog_upper = catalog.current_upper().await;
5223        boot_ts = std::cmp::max(boot_ts, catalog_upper);
5224
5225        if !read_only_controllers {
5226            epoch_millis_oracle.apply_write(boot_ts).await;
5227        }
5228
5229        info!(
5230            "startup: coordinator init: catalog open complete in {:?}",
5231            catalog_open_start.elapsed()
5232        );
5233
5234        // Whether replacement-migrated builtin MVs may write their new shards before cut-over.
5235        // Both `bootstrap` and the readiness gate below read this, and they have to agree.
5236        // `MIN_LEADER_VERSION_FOR_MIGRATED_MV_WRITES` explains why the leader's version settles it.
5237        //
5238        // While we are read-only, `last_seen_version` is that leader's version: our catalog
5239        // transaction is a savepoint, so our own bump of the setting never lands. `None` means a
5240        // freshly initialized catalog, with nothing migrated and no leader to be compatible with.
5241        //
5242        // `ENABLE_0DT_HYDRATE_MIGRATED_BUILTIN_MVS` is the break-glass revert: off falls back to
5243        // excluding migrated MVs from the caught-up gate, no redeploy needed.
5244        let hydrate_migrated_mvs = ENABLE_0DT_HYDRATE_MIGRATED_BUILTIN_MVS
5245            .get(catalog.system_config().dyncfgs())
5246            && last_seen_version
5247                .as_ref()
5248                .is_none_or(|version| *version >= MIN_LEADER_VERSION_FOR_MIGRATED_MV_WRITES);
5249
5250        let coord_thread_start = Instant::now();
5251        info!("startup: coordinator init: coordinator thread start beginning");
5252
5253        let session_id = catalog.config().session_id;
5254        let start_instant = catalog.config().start_instant;
5255
5256        // In order for the coordinator to support Rc and Refcell types, it cannot be
5257        // sent across threads. Spawn it in a thread and have this parent thread wait
5258        // for bootstrap completion before proceeding.
5259        let (bootstrap_tx, bootstrap_rx) = oneshot::channel();
5260        let handle = TokioHandle::current();
5261
5262        let metrics = Metrics::register_into(&metrics_registry);
5263        let metrics_clone = metrics.clone();
5264        let optimizer_metrics = OptimizerMetrics::register_into(
5265            &metrics_registry,
5266            catalog.system_config().optimizer_e2e_latency_warning_threshold(),
5267        );
5268        let segment_client_clone = segment_client.clone();
5269        let coord_now = now.clone();
5270        let advance_timelines_interval =
5271            tokio::time::interval(catalog.system_config().default_timestamp_interval());
5272
5273        let clusters_caught_up_check_interval = if read_only_controllers {
5274            let dyncfgs = catalog.system_config().dyncfgs();
5275            let interval = WITH_0DT_DEPLOYMENT_CAUGHT_UP_CHECK_INTERVAL.get(dyncfgs);
5276
5277            let mut interval = tokio::time::interval(interval);
5278            interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
5279            interval
5280        } else {
5281            // When not in read-only mode, we don't do hydration checks. But we
5282            // still have to provide _some_ interval. This is large enough that
5283            // it doesn't matter.
5284            //
5285            // TODO(aljoscha): We cannot use Duration::MAX right now because of
5286            // https://github.com/tokio-rs/tokio/issues/6634. Use that once it's
5287            // fixed for good.
5288            let mut interval = tokio::time::interval(Duration::from_secs(60 * 60));
5289            interval.set_missed_tick_behavior(MissedTickBehavior::Skip);
5290            interval
5291        };
5292
5293        let clusters_caught_up_check =
5294            clusters_caught_up_trigger.map(|trigger| {
5295                let mut exclude_collections: BTreeSet<GlobalId> =
5296                    new_builtin_collections.iter().copied().collect();
5297
5298                // A collection that can't advance its write frontier in read-only mode
5299                // stalls its transitive dependents too, so exclude those from the caught-up
5300                // check as well. That's every *new* builtin collection, whose fresh shard has no
5301                // writer until this deployment promotes, plus migrated MVs whenever the leader is
5302                // too old for them to write. An excluded dependent may still be hydrating right
5303                // after promotion, a brief blip we accept because these collections are small and
5304                // get a writer at cut-over.
5305                //
5306                // Seeded from all of `new_builtin_collections`, not just the MVs: a new builtin
5307                // table or source has no read-only writer either (`register_table_collections`
5308                // retains only *migrated* tables), so an MV reading one never advances past its
5309                // empty frontier. A *migrated* table is the opposite case, even though a builtin
5310                // MV can read one (`mz_clusters` joins `mz_cluster_replica_size_internal`):
5311                // `read_only_mode_table_worker` keeps advancing migrated tables' uppers.
5312                let new_builtin_items = new_builtin_collections.iter().map(|global_id| {
5313                    catalog
5314                        .state()
5315                        .try_get_entry_by_global_id(global_id)
5316                        .expect("new builtin collections have catalog entries")
5317                        .id()
5318                });
5319                let frozen_migrated_mvs = migrated_storage_collections_0dt
5320                    .iter()
5321                    .copied()
5322                    .filter(|_| !hydrate_migrated_mvs)
5323                    .filter(|id| catalog.state().get_entry(id).is_materialized_view());
5324                let mut todo: Vec<_> = new_builtin_items.chain(frozen_migrated_mvs).collect();
5325                while let Some(item_id) = todo.pop() {
5326                    let entry = catalog.state().get_entry(&item_id);
5327                    exclude_collections.extend(entry.global_ids());
5328                    todo.extend_from_slice(entry.used_by());
5329                }
5330
5331                CaughtUpCheckContext {
5332                    trigger,
5333                    exclude_collections,
5334                    cluster_stability: BTreeMap::new(),
5335                }
5336            });
5337
5338        if let Some(TimestampOracleConfig::Postgres(pg_config)) =
5339            timestamp_oracle_config.as_ref()
5340        {
5341            // Apply settings from system vars as early as possible because some
5342            // of them are locked in right when an oracle is first opened!
5343            let pg_timestamp_oracle_params =
5344                flags::timestamp_oracle_config(catalog.system_config());
5345            pg_timestamp_oracle_params.apply(pg_config);
5346        }
5347
5348        // Register a callback so whenever the MAX_CONNECTIONS or SUPERUSER_RESERVED_CONNECTIONS
5349        // system variables change, we update our connection limits.
5350        let connection_limit_callback: Arc<dyn Fn(&SystemVars) + Send + Sync> =
5351            Arc::new(move |system_vars: &SystemVars| {
5352                let limit: u64 = system_vars.max_connections().cast_into();
5353                let superuser_reserved: u64 =
5354                    system_vars.superuser_reserved_connections().cast_into();
5355
5356                // If superuser_reserved > max_connections, prefer max_connections.
5357                //
5358                // In this scenario all normal users would be locked out because all connections
5359                // would be reserved for superusers so complain if this is the case.
5360                let superuser_reserved = if superuser_reserved >= limit {
5361                    tracing::warn!(
5362                        "superuser_reserved ({superuser_reserved}) is greater than max connections ({limit})!"
5363                    );
5364                    limit
5365                } else {
5366                    superuser_reserved
5367                };
5368
5369                (connection_limit_callback)(limit, superuser_reserved);
5370            });
5371        catalog.system_config_mut().register_callback(
5372            &mz_sql::session::vars::MAX_CONNECTIONS,
5373            Arc::clone(&connection_limit_callback),
5374        );
5375        catalog.system_config_mut().register_callback(
5376            &mz_sql::session::vars::SUPERUSER_RESERVED_CONNECTIONS,
5377            connection_limit_callback,
5378        );
5379
5380        let (group_commit_tx, group_commit_rx) = appends::notifier();
5381
5382        let parent_span = tracing::Span::current();
5383        let thread = thread::Builder::new()
5384            // The Coordinator thread tends to keep a lot of data on its stack. To
5385            // prevent a stack overflow we allocate a stack three times as big as the default
5386            // stack.
5387            .stack_size(3 * stack::STACK_SIZE)
5388            .name("coordinator".to_string())
5389            .spawn(move || {
5390                let span = info_span!(parent: parent_span, "coord::coordinator").entered();
5391
5392                let controller = handle
5393                    .block_on({
5394                        catalog.initialize_controller(
5395                            controller_config,
5396                            controller_envd_epoch,
5397                            read_only_controllers,
5398                        )
5399                    })
5400                    .unwrap_or_terminate("failed to initialize storage_controller");
5401                // Initializing the controller uses one or more timestamps, so push the boot timestamp up to the
5402                // current catalog upper.
5403                let catalog_upper = handle.block_on(catalog.current_upper());
5404                boot_ts = std::cmp::max(boot_ts, catalog_upper);
5405                if !read_only_controllers {
5406                    let epoch_millis_oracle = &timestamp_oracles
5407                        .get(&Timeline::EpochMilliseconds)
5408                        .expect("inserted above")
5409                        .oracle;
5410                    handle.block_on(epoch_millis_oracle.apply_write(boot_ts));
5411                }
5412
5413                let catalog = Arc::new(catalog);
5414                // Both are read once at startup, see the field docs on
5415                // `occ_write_semaphore` and `frontend_read_then_write_enabled`.
5416                let max_concurrent_occ_writes =
5417                    usize::cast_from(catalog.system_config().max_concurrent_occ_writes());
5418                let frontend_read_then_write_enabled = {
5419                                FRONTEND_READ_THEN_WRITE.get(catalog.system_config().dyncfgs())
5420                };
5421
5422                let caching_secrets_reader = CachingSecretsReader::new(secrets_controller.reader());
5423                let (group_committer_tx, group_committer_rx) = mpsc::unbounded_channel();
5424                let mut coord = Coordinator {
5425                    controller,
5426                    catalog,
5427                    internal_cmd_tx,
5428                    group_commit_tx,
5429                    reconcile_now: Arc::new(Notify::new()),
5430                    group_committer_tx,
5431                    strict_serializable_reads_tx,
5432                    linearize_reads_notify: Arc::new(Notify::new()),
5433                    global_timelines: timestamp_oracles,
5434                    transient_id_gen: Arc::new(TransientIdGen::new()),
5435                    active_conns: BTreeMap::new(),
5436                    txn_read_holds: Default::default(),
5437                    pending_peeks: BTreeMap::new(),
5438                    client_pending_peeks: BTreeMap::new(),
5439                    pending_linearize_read_txns: BTreeMap::new(),
5440                    serialized_ddl: LockedVecDeque::new(),
5441                    active_compute_sinks: BTreeMap::new(),
5442                    active_webhooks: BTreeMap::new(),
5443                    active_copies: BTreeMap::new(),
5444                    connection_cancel_watches: BTreeMap::new(),
5445                    introspection_subscribes: BTreeMap::new(),
5446                    hydration_history_replica_cursor: None,
5447                    hydration_history_sweep: None,
5448                    metric_sinks: BTreeMap::new(),
5449                    metric_sink_plans: BTreeMap::new(),
5450                    write_locks: BTreeMap::new(),
5451                    deferred_write_ops: BTreeMap::new(),
5452                    pending_writes: Vec::new(),
5453                    occ_write_semaphore: Arc::new(Semaphore::new(max_concurrent_occ_writes)),
5454                    frontend_read_then_write_enabled,
5455                    advance_timelines_interval,
5456                    secrets_controller,
5457                    caching_secrets_reader,
5458                    cloud_resource_controller,
5459                    storage_usage_client,
5460                    storage_usage_collection_interval,
5461                    segment_client,
5462                    metrics,
5463                    catalog_info_metrics_registry: metrics_registry.clone(),
5464                    scoped_frontend: None,
5465                    optimizer_metrics,
5466                    tracing_handle,
5467                    statement_logging: StatementLogging::new(coord_now.clone()),
5468                    webhook_concurrency_limit,
5469                    timestamp_oracle_config,
5470                    caught_up_check_interval: clusters_caught_up_check_interval,
5471                    caught_up_check: clusters_caught_up_check,
5472                    installed_watch_sets: BTreeMap::new(),
5473                    connection_watch_sets: BTreeMap::new(),
5474                    cluster_replica_statuses: ClusterReplicaStatuses::new(),
5475                    read_only_controllers,
5476                    buffered_builtin_table_updates: Some(Vec::new()),
5477                    license_key,
5478                    user_id_pool: IdPool::empty(),
5479                    persist_client,
5480                };
5481
5482                // Read-only promotion restarts the process and creates a fresh committer.
5483                handle.block_on(async {
5484                    appends::spawn_group_committer(
5485                        group_committer_rx,
5486                        coord.get_local_timestamp_oracle(),
5487                        coord.controller.storage.table_write_handle(),
5488                        coord.catalog().upper_handle(),
5489                        coord.internal_cmd_tx.clone(),
5490                        coord.catalog().config().now.clone(),
5491                        coord.metrics.clone(),
5492                        coord.catalog().system_config().dyncfgs(),
5493                    );
5494                });
5495
5496                let bootstrap = handle.block_on(async {
5497                    coord
5498                        .bootstrap(
5499                            boot_ts,
5500                            migrated_storage_collections_0dt,
5501                            hydrate_migrated_mvs,
5502                            builtin_table_updates,
5503                            cached_global_exprs,
5504                            uncached_local_exprs,
5505                        )
5506                        .await?;
5507                    coord
5508                        .controller
5509                        .remove_orphaned_replicas(
5510                            coord.catalog().get_next_user_replica_id().await?,
5511                            coord.catalog().get_next_system_replica_id().await?,
5512                        )
5513                        .await
5514                        .map_err(AdapterError::Orchestrator)?;
5515
5516                    if let Some(retention_period) = storage_usage_retention_period {
5517                        coord
5518                            .prune_storage_usage_events_on_startup(retention_period)
5519                            .await;
5520                    }
5521
5522                    coord.prune_arrangement_sizes_history_on_startup().await;
5523
5524                    Ok(())
5525                });
5526                let ok = bootstrap.is_ok();
5527                drop(span);
5528                bootstrap_tx
5529                    .send(bootstrap)
5530                    .expect("bootstrap_rx is not dropped until it receives this message");
5531                if ok {
5532                    handle.block_on(coord.serve(
5533                        internal_cmd_rx,
5534                        strict_serializable_reads_rx,
5535                        cmd_rx,
5536                        group_commit_rx,
5537                    ));
5538                }
5539            })
5540            .expect("failed to create coordinator thread");
5541        match bootstrap_rx
5542            .await
5543            .expect("bootstrap_tx always sends a message or panics/halts")
5544        {
5545            Ok(()) => {
5546                info!(
5547                    "startup: coordinator init: coordinator thread start complete in {:?}",
5548                    coord_thread_start.elapsed()
5549                );
5550                info!(
5551                    "startup: coordinator init: complete in {:?}",
5552                    coord_start.elapsed()
5553                );
5554                let handle = Handle {
5555                    session_id,
5556                    start_instant,
5557                    _thread: thread.join_on_drop(),
5558                };
5559                let client = Client::new(
5560                    build_info,
5561                    cmd_tx,
5562                    metrics_clone,
5563                    now,
5564                    environment_id,
5565                    segment_client_clone,
5566                );
5567                Ok((handle, client))
5568            }
5569            Err(e) => Err(e),
5570        }
5571    }
5572    .boxed()
5573}
5574
5575// Determines and returns the highest timestamp for each timeline, for all known
5576// timestamp oracle implementations.
5577//
5578// Initially, we did this so that we can switch between implementations of
5579// timestamp oracle, but now we also do this to determine a monotonic boot
5580// timestamp, a timestamp that does not regress across reboots.
5581//
5582// This mostly works, but there can be linearizability violations, because there
5583// is no central moment where we do distributed coordination for all oracle
5584// types. Working around this seems prohibitively hard, maybe even impossible so
5585// we have to live with this window of potential violations during the upgrade
5586// window (which is the only point where we should switch oracle
5587// implementations).
5588async fn get_initial_oracle_timestamps(
5589    timestamp_oracle_config: &Option<TimestampOracleConfig>,
5590) -> Result<BTreeMap<Timeline, Timestamp>, AdapterError> {
5591    let mut initial_timestamps = BTreeMap::new();
5592
5593    if let Some(config) = timestamp_oracle_config {
5594        let oracle_timestamps = config.get_all_timelines().await?;
5595
5596        let debug_msg = || {
5597            oracle_timestamps
5598                .iter()
5599                .map(|(timeline, ts)| format!("{:?} -> {}", timeline, ts))
5600                .join(", ")
5601        };
5602        info!(
5603            "current timestamps from the timestamp oracle: {}",
5604            debug_msg()
5605        );
5606
5607        for (timeline, ts) in oracle_timestamps {
5608            let entry = initial_timestamps
5609                .entry(Timeline::from_str(&timeline).expect("could not parse timeline"));
5610
5611            entry
5612                .and_modify(|current_ts| *current_ts = std::cmp::max(*current_ts, ts))
5613                .or_insert(ts);
5614        }
5615    } else {
5616        info!("no timestamp oracle configured!");
5617    };
5618
5619    let debug_msg = || {
5620        initial_timestamps
5621            .iter()
5622            .map(|(timeline, ts)| format!("{:?}: {}", timeline, ts))
5623            .join(", ")
5624    };
5625    info!("initial oracle timestamps: {}", debug_msg());
5626
5627    Ok(initial_timestamps)
5628}
5629
5630#[instrument]
5631pub async fn load_remote_system_parameters(
5632    storage: &mut Box<dyn OpenableDurableCatalogState>,
5633    system_parameter_sync_config: Option<SystemParameterSyncConfig>,
5634    system_parameter_sync_timeout: Duration,
5635) -> Result<Option<BTreeMap<String, String>>, AdapterError> {
5636    if let Some(system_parameter_sync_config) = system_parameter_sync_config {
5637        tracing::info!("parameter sync on boot: start sync");
5638
5639        // We intentionally block initial startup, potentially forever,
5640        // on initializing LaunchDarkly. This may seem scary, but the
5641        // alternative is even scarier. Over time, we expect that the
5642        // compiled-in default values for the system parameters will
5643        // drift substantially from the defaults configured in
5644        // LaunchDarkly, to the point that starting an environment
5645        // without loading the latest values from LaunchDarkly will
5646        // result in running an untested configuration.
5647        //
5648        // Note this only applies during initial startup. Restarting
5649        // after we've synced once only blocks for a maximum of
5650        // `FRONTEND_SYNC_TIMEOUT` on LaunchDarkly, as it seems
5651        // reasonable to assume that the last-synced configuration was
5652        // valid enough.
5653        //
5654        // This philosophy appears to provide a good balance between not
5655        // running untested configurations in production while also not
5656        // making LaunchDarkly a "tier 1" dependency for existing
5657        // environments.
5658        //
5659        // If this proves to be an issue, we could seek to address the
5660        // configuration drift in a different way--for example, by
5661        // writing a script that runs in CI nightly and checks for
5662        // deviation between the compiled Rust code and LaunchDarkly.
5663        //
5664        // If it is absolutely necessary to bring up a new environment
5665        // while LaunchDarkly is down, the following manual mitigation
5666        // can be performed:
5667        //
5668        //    1. Edit the environmentd startup parameters to omit the
5669        //       LaunchDarkly configuration.
5670        //    2. Boot environmentd.
5671        //    3. Use the catalog-debug tool to run `edit config "{\"key\":\"system_config_synced\"}" "{\"value\": 1}"`.
5672        //    4. Adjust any other parameters as necessary to avoid
5673        //       running a nonstandard configuration in production.
5674        //    5. Edit the environmentd startup parameters to restore the
5675        //       LaunchDarkly configuration, for when LaunchDarkly comes
5676        //       back online.
5677        //    6. Reboot environmentd.
5678        let mut params = SynchronizedParameters::new(SystemVars::default());
5679        let frontend_sync = async {
5680            let frontend = SystemParameterFrontend::from(&system_parameter_sync_config).await?;
5681            frontend.pull(&mut params);
5682            let ops = params
5683                .modified()
5684                .into_iter()
5685                .map(|param| {
5686                    let name = param.name;
5687                    let value = param.value;
5688                    tracing::info!(name, value, initial = true, "sync parameter");
5689                    (name, value)
5690                })
5691                .collect();
5692            tracing::info!("parameter sync on boot: end sync");
5693            Ok(Some(ops))
5694        };
5695        if !storage.has_system_config_synced_once().await? {
5696            frontend_sync.await
5697        } else {
5698            match mz_ore::future::timeout(system_parameter_sync_timeout, frontend_sync).await {
5699                Ok(ops) => Ok(ops),
5700                Err(TimeoutError::Inner(e)) => Err(e),
5701                Err(TimeoutError::DeadlineElapsed) => {
5702                    tracing::info!("parameter sync on boot: sync has timed out");
5703                    Ok(None)
5704                }
5705            }
5706        }
5707    } else {
5708        Ok(None)
5709    }
5710}
5711
5712#[derive(Debug)]
5713pub enum WatchSetResponse {
5714    StatementDependenciesReady(StatementLoggingId, StatementLifecycleEvent),
5715    AlterSinkReady(AlterSinkReadyContext),
5716    AlterMaterializedViewReady(AlterMaterializedViewReadyContext),
5717}
5718
5719#[derive(Debug)]
5720pub struct AlterSinkReadyContext {
5721    ctx: Option<ExecuteContext>,
5722    otel_ctx: OpenTelemetryContext,
5723    plan: AlterSinkPlan,
5724    plan_validity: PlanValidity,
5725    read_hold: ReadHolds,
5726}
5727
5728impl AlterSinkReadyContext {
5729    fn ctx(&mut self) -> &mut ExecuteContext {
5730        self.ctx.as_mut().expect("only cleared on drop")
5731    }
5732
5733    fn retire(mut self, result: Result<ExecuteResponse, AdapterError>) {
5734        self.ctx
5735            .take()
5736            .expect("only cleared on drop")
5737            .retire(result);
5738    }
5739}
5740
5741impl Drop for AlterSinkReadyContext {
5742    fn drop(&mut self) {
5743        if let Some(ctx) = self.ctx.take() {
5744            ctx.retire(Err(AdapterError::Canceled));
5745        }
5746    }
5747}
5748
5749#[derive(Debug)]
5750pub struct AlterMaterializedViewReadyContext {
5751    ctx: Option<ExecuteContext>,
5752    otel_ctx: OpenTelemetryContext,
5753    plan: plan::AlterMaterializedViewApplyReplacementPlan,
5754    plan_validity: PlanValidity,
5755}
5756
5757impl AlterMaterializedViewReadyContext {
5758    fn ctx(&mut self) -> &mut ExecuteContext {
5759        self.ctx.as_mut().expect("only cleared on drop")
5760    }
5761
5762    fn retire(mut self, result: Result<ExecuteResponse, AdapterError>) {
5763        self.ctx
5764            .take()
5765            .expect("only cleared on drop")
5766            .retire(result);
5767    }
5768}
5769
5770impl Drop for AlterMaterializedViewReadyContext {
5771    fn drop(&mut self) {
5772        if let Some(ctx) = self.ctx.take() {
5773            ctx.retire(Err(AdapterError::Canceled));
5774        }
5775    }
5776}
5777
5778/// A struct for tracking the ownership of a lock and a VecDeque to store to-be-done work after the
5779/// lock is freed.
5780#[derive(Debug)]
5781struct LockedVecDeque<T> {
5782    items: VecDeque<T>,
5783    lock: Arc<tokio::sync::Mutex<()>>,
5784}
5785
5786impl<T> LockedVecDeque<T> {
5787    pub fn new() -> Self {
5788        Self {
5789            items: VecDeque::new(),
5790            lock: Arc::new(tokio::sync::Mutex::new(())),
5791        }
5792    }
5793
5794    pub fn try_lock_owned(&self) -> Result<OwnedMutexGuard<()>, tokio::sync::TryLockError> {
5795        Arc::clone(&self.lock).try_lock_owned()
5796    }
5797
5798    pub fn is_empty(&self) -> bool {
5799        self.items.is_empty()
5800    }
5801
5802    pub fn push_back(&mut self, value: T) {
5803        self.items.push_back(value)
5804    }
5805
5806    pub fn pop_front(&mut self) -> Option<T> {
5807        self.items.pop_front()
5808    }
5809
5810    pub fn remove(&mut self, index: usize) -> Option<T> {
5811        self.items.remove(index)
5812    }
5813
5814    pub fn iter(&self) -> std::collections::vec_deque::Iter<'_, T> {
5815        self.items.iter()
5816    }
5817}
5818
5819#[derive(Debug)]
5820struct DeferredPlanStatement {
5821    ctx: ExecuteContext,
5822    ps: PlanStatement,
5823}
5824
5825#[derive(Debug)]
5826enum PlanStatement {
5827    Statement {
5828        stmt: Arc<Statement<Raw>>,
5829        params: Params,
5830    },
5831    Plan {
5832        plan: mz_sql::plan::Plan,
5833        resolved_ids: ResolvedIds,
5834        sql_impl_resolved_ids: ResolvedIds,
5835    },
5836}
5837
5838#[derive(Debug, Error)]
5839pub enum NetworkPolicyError {
5840    #[error("Access denied for address {0}")]
5841    AddressDenied(IpAddr),
5842    #[error("Access denied missing IP address")]
5843    MissingIp,
5844}
5845
5846pub(crate) fn validate_ip_with_policy_rules(
5847    ip: &IpAddr,
5848    rules: &Vec<NetworkPolicyRule>,
5849) -> Result<(), NetworkPolicyError> {
5850    // At the moment we're not handling action or direction
5851    // as those are only able to be "allow" and "ingress" respectively
5852    if rules.iter().any(|r| r.address.0.contains(ip)) {
5853        Ok(())
5854    } else {
5855        Err(NetworkPolicyError::AddressDenied(ip.clone()))
5856    }
5857}
5858
5859pub(crate) fn infer_sql_type_for_catalog(
5860    hir_expr: &HirRelationExpr,
5861    mir_expr: &MirRelationExpr,
5862) -> SqlRelationType {
5863    let mut typ = hir_expr.top_level_typ();
5864    typ.backport_nullability_and_keys(&mir_expr.typ());
5865    typ
5866}
5867
5868#[cfg(test)]
5869mod execute_context_tests {
5870    use tokio::sync::{mpsc, oneshot};
5871
5872    use super::*;
5873    use crate::session::Session;
5874    use crate::util::ClientTransmitter;
5875
5876    /// Runtime shutdown drops the barrier-waiting task that `retire` spawns. The context's `Drop`
5877    /// backstop must answer the client, rather than panicking on an unsent `ClientTransmitter`.
5878    #[mz_ore::test]
5879    fn test_retire_answers_client_when_runtime_shuts_down() {
5880        let runtime = tokio::runtime::Runtime::new().expect("can build runtime");
5881
5882        let (client_tx, mut client_rx) = oneshot::channel();
5883        let (internal_cmd_tx, _internal_cmd_rx) = mpsc::unbounded_channel();
5884
5885        runtime.block_on(async {
5886            let ctx = ExecuteContext::from_parts_with_response_barriers(
5887                ClientTransmitter::new(client_tx, internal_cmd_tx.clone()),
5888                internal_cmd_tx,
5889                Session::dummy(),
5890                ExecuteContextGuard::default(),
5891                // Stands in for a group commit that shutdown will never apply.
5892                vec![Box::pin(std::future::pending())],
5893            );
5894            ctx.retire(Ok(ExecuteResponse::StartedTransaction));
5895        });
5896
5897        drop(runtime);
5898
5899        let response = client_rx.try_recv().expect("client must be answered");
5900        assert!(
5901            matches!(response.result, Err(AdapterError::Internal(_))),
5902            "expected an internal error, got {:?}",
5903            response.result
5904        );
5905    }
5906}
5907
5908#[cfg(test)]
5909mod id_pool_tests {
5910    use super::IdPool;
5911
5912    #[mz_ore::test]
5913    fn test_empty_pool() {
5914        let mut pool = IdPool::empty();
5915        assert_eq!(pool.remaining(), 0);
5916        assert_eq!(pool.allocate(), None);
5917        assert_eq!(pool.allocate_many(1), None);
5918    }
5919
5920    #[mz_ore::test]
5921    fn test_allocate_single() {
5922        let mut pool = IdPool::empty();
5923        pool.refill(10, 13);
5924        assert_eq!(pool.remaining(), 3);
5925        assert_eq!(pool.allocate(), Some(10));
5926        assert_eq!(pool.allocate(), Some(11));
5927        assert_eq!(pool.allocate(), Some(12));
5928        assert_eq!(pool.remaining(), 0);
5929        assert_eq!(pool.allocate(), None);
5930    }
5931
5932    #[mz_ore::test]
5933    fn test_allocate_many() {
5934        let mut pool = IdPool::empty();
5935        pool.refill(100, 105);
5936        assert_eq!(pool.allocate_many(3), Some(vec![100, 101, 102]));
5937        assert_eq!(pool.remaining(), 2);
5938        // Not enough remaining for 3 more.
5939        assert_eq!(pool.allocate_many(3), None);
5940        // But 2 works.
5941        assert_eq!(pool.allocate_many(2), Some(vec![103, 104]));
5942        assert_eq!(pool.remaining(), 0);
5943    }
5944
5945    #[mz_ore::test]
5946    fn test_allocate_many_zero() {
5947        let mut pool = IdPool::empty();
5948        pool.refill(1, 5);
5949        assert_eq!(pool.allocate_many(0), Some(vec![]));
5950        assert_eq!(pool.remaining(), 4);
5951    }
5952
5953    #[mz_ore::test]
5954    fn test_refill_resets_pool() {
5955        let mut pool = IdPool::empty();
5956        pool.refill(0, 2);
5957        assert_eq!(pool.allocate(), Some(0));
5958        // Refill before exhaustion replaces the range.
5959        pool.refill(50, 52);
5960        assert_eq!(pool.allocate(), Some(50));
5961        assert_eq!(pool.allocate(), Some(51));
5962        assert_eq!(pool.allocate(), None);
5963    }
5964
5965    #[mz_ore::test]
5966    fn test_mixed_allocate_and_allocate_many() {
5967        let mut pool = IdPool::empty();
5968        pool.refill(0, 10);
5969        assert_eq!(pool.allocate(), Some(0));
5970        assert_eq!(pool.allocate_many(3), Some(vec![1, 2, 3]));
5971        assert_eq!(pool.allocate(), Some(4));
5972        assert_eq!(pool.remaining(), 5);
5973    }
5974
5975    #[mz_ore::test]
5976    #[should_panic(expected = "invalid pool range")]
5977    fn test_refill_invalid_range_panics() {
5978        let mut pool = IdPool::empty();
5979        pool.refill(10, 5);
5980    }
5981}
5982
5983#[cfg(test)]
5984mod arrangement_sizes_pruner_tests {
5985    use mz_repr::catalog_item_id::CatalogItemId;
5986    use mz_repr::{Datum, Row};
5987
5988    use super::arrangement_sizes_expired_retractions;
5989
5990    // Pack a row shaped like `mz_object_arrangement_size_history`: the pruner
5991    // only cares about column 3 (`collection_timestamp`), but we stuff the
5992    // other three columns with realistic values so shape changes would fail.
5993    fn history_row(ts_ms: i64) -> Row {
5994        let dt = mz_ore::now::to_datetime(ts_ms.try_into().expect("non-negative"));
5995        Row::pack_slice(&[
5996            Datum::String("r1"),
5997            Datum::String("u1"),
5998            Datum::Int64(123),
5999            Datum::TimestampTz(dt.try_into().expect("fits in TimestampTz")),
6000        ])
6001    }
6002
6003    fn item_id() -> CatalogItemId {
6004        // Any CatalogItemId will do; tests don't dispatch on it.
6005        CatalogItemId::User(42)
6006    }
6007
6008    #[mz_ore::test]
6009    fn empty_input_produces_no_retractions() {
6010        let out = arrangement_sizes_expired_retractions(Vec::new(), 1_000, item_id());
6011        assert!(out.is_empty());
6012    }
6013
6014    #[mz_ore::test]
6015    fn retracts_only_rows_strictly_before_cutoff() {
6016        // Mixes both sides of the filter and includes a row at exactly
6017        // the cutoff timestamp to pin down the strict-less-than boundary.
6018        let rows = vec![
6019            (history_row(100), 1),
6020            (history_row(500), 1),
6021            (history_row(1_000), 1), // at cutoff: kept (strict <)
6022            (history_row(5_000), 1),
6023        ];
6024        let out = arrangement_sizes_expired_retractions(rows, 1_000, item_id());
6025        assert_eq!(out.len(), 2);
6026    }
6027
6028    #[mz_ore::test]
6029    #[should_panic(expected = "consolidated contents should not contain retractions")]
6030    fn retraction_in_input_panics() {
6031        let rows = vec![(history_row(100), -1)];
6032        let _ = arrangement_sizes_expired_retractions(rows, 1_000, item_id());
6033    }
6034}