1use 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
255const MIN_LEADER_VERSION_FOR_MIGRATED_MV_WRITES: Version = Version::new(26, 17, 0);
265
266#[derive(Debug)]
292pub(crate) struct IdPool {
293 next: u64,
294 upper: u64,
295}
296
297impl IdPool {
298 pub fn empty() -> Self {
300 IdPool { next: 0, upper: 0 }
301 }
302
303 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 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 pub fn remaining(&self) -> u64 {
328 self.upper - self.next
329 }
330
331 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#[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 conn_id: ConnectionId,
362 acquired_lock: Option<(CatalogItemId, tokio::sync::OwnedMutexGuard<()>)>,
372 },
373 GroupCommitInitiate(Span, Option<GroupCommitPermit>),
375 GroupCommitApplied {
379 responses: Vec<crate::util::CompletedClientTransmitter>,
381 statement_logging_ids: Vec<StatementLoggingId>,
383 internal_results: Vec<crate::coord::appends::InternalWriteResponder>,
385 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 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 ClusterControllerRequest(cluster_controller::ClusterControllerRequest),
483}
484
485impl Message {
486 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#[derive(Debug)]
601pub enum ControllerReadiness {
602 Storage,
604 Compute,
606 Metrics,
608 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 LinearizeTimestamp(PeekStageLinearizeTimestamp),
646 RealTimeRecency(PeekStageRealTimeRecency),
647 TimestampReadHold(PeekStageTimestampReadHold),
648 Optimize(PeekStageOptimize),
649 Finish(PeekStageFinish),
651 ExplainPlan(PeekStageExplainPlan),
653 ExplainPushdown(PeekStageExplainPushdown),
654 CopyToPreflight(PeekStageCopyTo),
656 CopyToDataflow(PeekStageCopyTo),
658}
659
660#[derive(Debug)]
661pub struct CopyToContext {
662 pub desc: RelationDesc,
664 pub uri: Uri,
666 pub connection: StorageConnection<ReferencedConnection>,
668 pub connection_id: CatalogItemId,
670 pub format: S3SinkFormat,
672 pub max_file_size: u64,
674 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 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 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 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 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 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 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 explain_ctx: ExplainContext,
864}
865
866#[derive(Debug)]
867pub struct CreateViewFinish {
868 validity: PlanValidity,
869 item_id: CatalogItemId,
871 global_id: GlobalId,
873 plan: plan::CreateViewPlan,
874 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 timeline_context: TimelineContext,
933 oracle_read_ts: Option<Timestamp>,
937}
938
939#[derive(Debug)]
940pub enum ClusterStage {
941 Alter(AlterCluster),
942 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 target: ReconfigurationTarget,
963}
964
965#[derive(Debug)]
966pub enum ExplainContext {
967 None,
969 Plan(ExplainPlanContext),
971 PlanInsightsNotice(OptimizerTrace),
974 Pushdown,
976}
977
978impl ExplainContext {
979 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 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 explain_ctx: ExplainContext,
1041}
1042
1043#[derive(Debug)]
1044pub struct CreateMaterializedViewFinish {
1045 item_id: CatalogItemId,
1047 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 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 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 oracle_read_ts: Option<Timestamp>,
1116 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 sink_id: GlobalId,
1187 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#[derive(Debug, Copy, Clone, PartialEq, Eq)]
1250pub enum TargetCluster {
1251 CatalogServer,
1253 Active,
1255 Transaction(ClusterId),
1257}
1258
1259pub(crate) enum StageResult<T> {
1261 Handle(JoinHandle<Result<T, AdapterError>>),
1263 HandleRetire(JoinHandle<Result<ExecuteResponse, AdapterError>>),
1265 Immediate(T),
1267 Response(ExecuteResponse),
1269}
1270
1271pub(crate) trait Staged: Send {
1273 type Ctx: StagedContext;
1274
1275 fn validity(&mut self) -> &mut PlanValidity;
1276
1277 async fn stage(
1279 self,
1280 coord: &mut Coordinator,
1281 ctx: &mut Self::Ctx,
1282 ) -> Result<StageResult<Box<Self>>, AdapterError>;
1283
1284 fn message(self, ctx: Self::Ctx, span: Span) -> Message;
1286
1287 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
1314pub 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 pub read_only_controllers: bool,
1353
1354 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#[derive(Debug, Serialize)]
1367pub struct ConnMeta {
1368 secret_key: u32,
1373 connected_at: EpochMillis,
1375 user: User,
1376 application_name: String,
1377 uuid: Uuid,
1378 conn_id: ConnectionId,
1379 client_ip: Option<IpAddr>,
1380
1381 drop_sinks: BTreeSet<GlobalId>,
1384
1385 #[serde(skip)]
1387 deferred_lock: Option<OwnedMutexGuard<()>>,
1388
1389 #[serde(skip)]
1391 notice_tx: mpsc::UnboundedSender<AdapterNotice>,
1392
1393 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)]
1430pub struct PendingTxn {
1432 ctx: ExecuteContext,
1434 response: Result<PendingTxnResponse, AdapterError>,
1436 action: EndTransactionAction,
1438}
1439
1440#[derive(Debug)]
1441pub enum PendingTxnResponse {
1443 Committed {
1445 params: BTreeMap<&'static str, String>,
1447 },
1448 Rolledback {
1450 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)]
1478pub struct PendingReadTxn {
1480 txn: PendingRead,
1482 timestamp_context: TimestampContext,
1484 created: Instant,
1486 num_requeues: u64,
1490 otel_ctx: OpenTelemetryContext,
1492}
1493
1494impl PendingReadTxn {
1495 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)]
1506enum PendingRead {
1508 Read {
1509 txn: PendingTxn,
1511 },
1512 ReadThenWrite {
1513 ctx: ExecuteContext,
1515 tx: oneshot::Sender<Option<ExecuteContext>>,
1518 },
1519}
1520
1521impl PendingRead {
1522 #[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 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 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 let _ = tx.send(None);
1569 ctx
1570 }
1571 }
1572 }
1573}
1574
1575#[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 #[must_use]
1604 pub(crate) fn retire(self) -> Option<StatementLoggingId> {
1605 self.statement_uuid
1606 }
1607}
1608
1609#[derive(Debug)]
1619#[must_use]
1620pub struct ExecuteContextGuard {
1621 extra: ExecuteContextExtra,
1622 coordinator_tx: mpsc::UnboundedSender<Message>,
1627}
1628
1629impl Default for ExecuteContextGuard {
1630 fn default() -> Self {
1631 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 pub(crate) fn defuse(mut self) -> ExecuteContextExtra {
1665 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 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 let _ = self.coordinator_tx.send(msg);
1685 }
1686 }
1687}
1688
1689#[derive(Debug)]
1694pub struct ExecuteContext {
1695 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 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 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 #[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 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 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 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 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 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 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 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 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 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 ClusterStatus::Offline(reason_x.or(reason_y))
2013 }
2014 })
2015 }
2016
2017 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 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 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#[derive(Derivative)]
2050#[derivative(Debug)]
2051pub struct Coordinator {
2052 #[derivative(Debug = "ignore")]
2054 controller: mz_controller::Controller,
2055 catalog: Arc<Catalog>,
2063
2064 persist_client: PersistClient,
2067
2068 internal_cmd_tx: mpsc::UnboundedSender<Message>,
2070 group_commit_tx: appends::GroupCommitNotifier,
2072 reconcile_now: Arc<Notify>,
2076 group_committer_tx: mpsc::UnboundedSender<appends::TableWriteCmd>,
2077
2078 strict_serializable_reads_tx: mpsc::UnboundedSender<(ConnectionId, PendingReadTxn)>,
2080
2081 linearize_reads_notify: Arc<Notify>,
2085
2086 global_timelines: BTreeMap<Timeline, TimelineState>,
2089
2090 transient_id_gen: Arc<TransientIdGen>,
2092 active_conns: BTreeMap<ConnectionId, ConnMeta>,
2095
2096 txn_read_holds: BTreeMap<ConnectionId, read_policy::ReadHolds>,
2100
2101 pending_peeks: BTreeMap<Uuid, PendingPeek>,
2105 client_pending_peeks: BTreeMap<ConnectionId, BTreeMap<Uuid, ClusterId>>,
2107
2108 pending_linearize_read_txns: BTreeMap<ConnectionId, PendingReadTxn>,
2110
2111 active_compute_sinks: BTreeMap<GlobalId, ActiveComputeSink>,
2113 active_webhooks: BTreeMap<CatalogItemId, WebhookAppenderInvalidator>,
2115 active_copies: BTreeMap<ConnectionId, ActiveCopyFrom>,
2118
2119 connection_cancel_watches: BTreeMap<ConnectionId, (watch::Sender<bool>, watch::Receiver<bool>)>,
2130 introspection_subscribes: BTreeMap<GlobalId, IntrospectionSubscribe>,
2132 hydration_history_replica_cursor: Option<ReplicaId>,
2134 hydration_history_sweep: Option<AbortOnDropHandle<()>>,
2136 metric_sinks: BTreeMap<(ReplicaId, &'static str), InstalledMetricSink>,
2141 metric_sink_plans: BTreeMap<&'static str, PlannedMetricSink>,
2144
2145 write_locks: BTreeMap<CatalogItemId, Arc<tokio::sync::Mutex<()>>>,
2147 deferred_write_ops: BTreeMap<ConnectionId, DeferredOp>,
2149
2150 pending_writes: Vec<PendingWriteTxn>,
2152
2153 occ_write_semaphore: Arc<Semaphore>,
2164
2165 frontend_read_then_write_enabled: bool,
2170
2171 advance_timelines_interval: Interval,
2181
2182 serialized_ddl: LockedVecDeque<DeferredPlanStatement>,
2191
2192 secrets_controller: Arc<dyn SecretsController>,
2195 caching_secrets_reader: CachingSecretsReader,
2197
2198 cloud_resource_controller: Option<Arc<dyn CloudResourceController>>,
2201
2202 storage_usage_client: StorageUsageClient,
2204 storage_usage_collection_interval: Duration,
2206
2207 #[derivative(Debug = "ignore")]
2209 segment_client: Option<mz_segment::Client>,
2210
2211 metrics: Metrics,
2213 optimizer_metrics: OptimizerMetrics,
2215
2216 tracing_handle: TracingHandle,
2218
2219 statement_logging: StatementLogging,
2221
2222 webhook_concurrency_limit: WebhookConcurrencyLimiter,
2224
2225 timestamp_oracle_config: Option<TimestampOracleConfig>,
2228
2229 caught_up_check_interval: Interval,
2232
2233 caught_up_check: Option<CaughtUpCheckContext>,
2236
2237 catalog_info_metrics_registry: MetricsRegistry,
2240
2241 scoped_frontend: Option<Arc<SystemParameterFrontend>>,
2250
2251 installed_watch_sets: BTreeMap<WatchSetId, (ConnectionId, WatchSetResponse)>,
2253
2254 connection_watch_sets: BTreeMap<ConnectionId, BTreeSet<WatchSetId>>,
2256
2257 cluster_replica_statuses: ClusterReplicaStatuses,
2259
2260 read_only_controllers: bool,
2264
2265 buffered_builtin_table_updates: Option<Vec<BuiltinTableUpdate>>,
2273
2274 license_key: ValidatedLicenseKey,
2275
2276 user_id_pool: IdPool,
2278}
2279
2280impl Coordinator {
2281 pub(crate) async fn reconcile_scoped_system_parameters(
2300 &mut self,
2301 scoped: ScopedParameters,
2302 prune_scope: ScopedParametersScope,
2303 ) {
2304 if self.catalog().state().scoped_system_parameters() == &scoped {
2307 return;
2308 }
2309
2310 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 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 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(¶ms, &cluster_param_names, &clusters);
2434 }
2435 if !replica_param_names.is_empty() && !replicas.is_empty() {
2436 evaluated.replica =
2437 frontend.pull_replica_overrides(¶ms, &replica_param_names, &replicas);
2438 }
2439 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 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 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 pub(crate) fn push_replica_dyncfg_overrides(&mut self) {
2504 let instance_overrides = self.replica_dyncfg_overrides();
2505
2506 self.controller
2521 .update_replica_dyncfg_overrides(instance_overrides);
2522 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 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 #[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 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 mz_metrics::update_dyncfg(&system_config.dyncfg_updates());
2598
2599 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 let replica_dyncfg_overrides = self.replica_dyncfg_overrides();
2617 self.controller
2618 .update_replica_dyncfg_overrides(replica_dyncfg_overrides);
2619
2620 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 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 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 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 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 CatalogItem::Source(source) => {
2760 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 self.catalog().state().pack_optimizer_notices(
2814 &mut builtin_table_updates,
2815 df_meta.optimizer_notices.iter(),
2816 Diff::ONE,
2817 );
2818 }
2819
2820 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 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 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 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 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 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 self.catalog().state().pack_optimizer_notices(
2925 &mut builtin_table_updates,
2926 df_meta.optimizer_notices.iter(),
2927 Diff::ONE,
2928 );
2929 }
2930
2931 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 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 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 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 drop(dataflow_read_holds);
2990 for (cw, policies) in policies_to_set {
2992 self.initialize_read_policies(&policies, cw).await;
2993 }
2994
2995 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 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 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 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 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 {
3103 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 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 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 self.controller.initialization_complete();
3183
3184 self.bootstrap_introspection_subscribes().await;
3186
3187 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 #[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 struct TableMetadata<'a> {
3215 id: CatalogItemId,
3216 name: &'a QualifiedItemName,
3217 table: &'a Table,
3218 }
3219
3220 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 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 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 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 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 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 let mut snapshot_cursor = snapshot_fut
3297 .await
3298 .unwrap_or_terminate("cannot fail to snapshot");
3299
3300 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 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 builtin_updates_fut.await;
3331 }
3332
3333 #[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 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 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 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 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 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 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 | CatalogItem::MetricSink(_) => (),
3542 }
3543 }
3544
3545 let register_ts = if self.controller.read_only() {
3546 self.get_local_read_ts().await
3547 } else {
3548 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 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 let mut pending: BTreeMap<_, _> = collections.into_iter().collect();
3582
3583 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 .filter(|dep_id| *dep_id != item_id)
3594 .map(|dep_id| self.catalog.get_entry(&dep_id).latest_global_id())
3595 .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 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 for (gid, collection) in &mut ready {
3621 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 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 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 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 #[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 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 .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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 }
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 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 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 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 let mut interval = tokio::time::interval(Duration::from_secs(5));
4145 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
4149
4150 let mut coord_stuck = false;
4152
4153 loop {
4154 interval.tick().await;
4155
4156 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 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 if coord_stuck {
4176 tracing::info!("Coordinator became unstuck");
4177 }
4178 coord_stuck = false;
4179
4180 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 let warn_threshold = self
4200 .catalog()
4201 .system_config()
4202 .coord_slow_message_warn_threshold();
4203
4204 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 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 select! {
4226 biased;
4231
4232 _ = internal_cmd_rx.recv_many(&mut messages, MESSAGE_BATCH) => {},
4236 Some(event) = cluster_events.next() => {
4240 messages.push(Message::ClusterEvent(event))
4241 },
4242 () = self.controller.ready() => {
4246 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 permit = group_commit_rx.ready() => {
4261 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 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 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 _ = self.advance_timelines_interval.tick() => {
4317 if self.controller.read_only() {
4321 messages.push(Message::AdvanceTimelines);
4322 } else {
4323 self.group_commit_tx.notify();
4324 }
4325 },
4326 () = linearize_reads_notified.as_mut() => {
4337 linearize_reads_notified.set(linearize_reads_notify.notified());
4338 messages.push(Message::LinearizeReads);
4339 }
4340 _ = self.caught_up_check_interval.tick() => {
4344 self.maybe_check_caught_up().await;
4349
4350 continue;
4351 },
4352
4353 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 message_batch.observe(f64::cast_lossy(messages.len()));
4369
4370 for msg in messages.drain(..) {
4371 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 *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 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 if let Some(sweep) = self.hydration_history_sweep.take() {
4428 sweep.abort_and_wait().await;
4429 }
4430
4431 if let Some(catalog) = Arc::into_inner(self.catalog) {
4434 catalog.expire().await;
4435 }
4436 }
4437 .boxed_local()
4438 }
4439
4440 fn catalog(&self) -> &Catalog {
4442 &self.catalog
4443 }
4444
4445 fn owned_catalog(&self) -> Arc<Catalog> {
4448 Arc::clone(&self.catalog)
4449 }
4450
4451 fn optimizer_metrics(&self) -> OptimizerMetrics {
4454 self.optimizer_metrics.clone()
4455 }
4456
4457 fn catalog_mut(&mut self) -> &mut Catalog {
4459 Arc::make_mut(&mut self.catalog)
4467 }
4468
4469 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, 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 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 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 fn connection_context(&self) -> &ConnectionContext {
4534 self.controller.connection_context()
4535 }
4536
4537 fn secrets_reader(&self) -> &Arc<dyn SecretsReader> {
4539 &self.connection_context().secrets_reader
4540 }
4541
4542 #[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 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 #[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 pub fn instance_snapshot(
4597 &self,
4598 id: ComputeInstanceId,
4599 ) -> Result<ComputeInstanceSnapshot, InstanceMissing> {
4600 ComputeInstanceSnapshot::new(&self.controller, id)
4601 }
4602
4603 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 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 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 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 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 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 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 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 pub async fn dump(&self) -> Result<serde_json::Value, anyhow::Error> {
4720 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 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 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 let _ = internal_cmd_tx.send(Message::StorageUsagePrune(expired));
4829 });
4830 }
4831
4832 async fn prune_arrangement_sizes_history_on_startup(&self) {
4841 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 let _ = internal_cmd_tx.send(Message::ArrangementSizesPrune(expired));
4872 });
4873 }
4874
4875 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 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
4910fn 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 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
4964struct LastMessage {
4966 kind: &'static str,
4967 stmt: Option<Arc<Statement<Raw>>>,
4968}
4969
4970impl LastMessage {
4971 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 if std::thread::panicking() {
4993 eprintln!("Coordinator panicking, dumping last message\n{self:?}",);
4995 }
4996 }
4997}
4998
4999pub 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 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 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(×tamp_oracle_config).await?;
5099
5100 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 let catalog_upper = storage.current_upper().await;
5123 let epoch_millis_oracle = ×tamp_oracles
5129 .get(&Timeline::EpochMilliseconds)
5130 .expect("inserted above")
5131 .oracle;
5132
5133 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 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 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 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 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 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 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 let pg_timestamp_oracle_params =
5344 flags::timestamp_oracle_config(catalog.system_config());
5345 pg_timestamp_oracle_params.apply(pg_config);
5346 }
5347
5348 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 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 .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 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 = ×tamp_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 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 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
5575async 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 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#[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 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 #[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 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 assert_eq!(pool.allocate_many(3), None);
5940 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 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 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 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 let rows = vec![
6019 (history_row(100), 1),
6020 (history_row(500), 1),
6021 (history_row(1_000), 1), (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}