Skip to main content

mz_adapter/catalog/
builtin_table_updates.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10mod notice;
11
12use bytesize::ByteSize;
13use ipnet::IpNet;
14use mz_adapter_types::compaction::CompactionWindow;
15use mz_audit_log::VersionedStorageUsage;
16use mz_catalog::SYSTEM_CONN_ID;
17use mz_catalog::builtin::{
18    BuiltinTable, MZ_AGGREGATES, MZ_ARRAY_TYPES, MZ_BASE_TYPES, MZ_CLUSTER_REPLICA_SIZE_INTERNAL,
19    MZ_CLUSTER_REPLICA_SIZES, MZ_COLUMNS, MZ_EGRESS_IPS, MZ_FUNCTIONS,
20    MZ_HISTORY_RETENTION_STRATEGIES, MZ_INDEX_COLUMNS, MZ_LICENSE_KEYS, MZ_LIST_TYPES,
21    MZ_MAP_TYPES, MZ_MATERIALIZED_VIEW_REFRESH_STRATEGIES, MZ_OPERATORS, MZ_PSEUDO_TYPES,
22    MZ_REPLACEMENTS, MZ_ROLE_AUTH, MZ_SESSIONS, MZ_STORAGE_USAGE_BY_SHARD, MZ_SUBSCRIPTIONS,
23    MZ_TYPE_PG_METADATA, MZ_TYPES, MZ_WEBHOOKS_SOURCES,
24};
25use mz_catalog::memory::error::Error;
26use mz_catalog::memory::objects::{
27    CatalogItem, DataSourceDesc, Func, Index, MaterializedView, Table, TableDataSource, Type,
28};
29use mz_expr::MirScalarExpr;
30use mz_license_keys::ValidatedLicenseKey;
31use mz_orchestrator::{CpuLimit, DiskLimit, MemoryLimit};
32use mz_ore::cast::CastFrom;
33use mz_ore::collections::CollectionExt;
34use mz_persist_client::batch::ProtoBatch;
35use mz_repr::adt::array::ArrayDimension;
36use mz_repr::adt::interval::Interval;
37use mz_repr::adt::jsonb::Jsonb;
38use mz_repr::adt::mz_acl_item::PrivilegeMap;
39use mz_repr::refresh_schedule::RefreshEvery;
40use mz_repr::role_id::RoleId;
41use mz_repr::{
42    CatalogItemId, Datum, Diff, GlobalId, ReprColumnType, Row, RowPacker, SqlScalarType, Timestamp,
43};
44use mz_sql::ast::{CreateIndexStatement, Statement};
45use mz_sql::catalog::{CatalogType, TypeCategory};
46use mz_sql::func::FuncImplCatalogDetails;
47use mz_sql::names::SchemaSpecifier;
48use mz_sql_parser::ast::display::AstDisplay;
49use mz_storage_client::client::TableData;
50use smallvec::smallvec;
51use uuid::Uuid;
52
53// DO NOT add any more imports from `crate` outside of `crate::catalog`.
54use crate::active_compute_sink::ActiveSubscribe;
55use crate::catalog::CatalogState;
56use crate::coord::ConnMeta;
57
58/// An update to a built-in table.
59#[derive(Debug, Clone)]
60pub struct BuiltinTableUpdate<T = CatalogItemId> {
61    /// The reference of the table to update.
62    pub id: T,
63    /// The data to put into the table.
64    pub data: TableData,
65}
66
67impl<T> BuiltinTableUpdate<T> {
68    /// Create a [`BuiltinTableUpdate`] from a [`Row`].
69    pub fn row(id: T, row: Row, diff: Diff) -> BuiltinTableUpdate<T> {
70        BuiltinTableUpdate {
71            id,
72            data: TableData::Rows(vec![(row, diff)]),
73        }
74    }
75
76    pub fn batch(id: T, batch: ProtoBatch) -> BuiltinTableUpdate<T> {
77        BuiltinTableUpdate {
78            id,
79            data: TableData::Batches(smallvec![batch]),
80        }
81    }
82}
83
84impl CatalogState {
85    pub fn resolve_builtin_table_updates(
86        &self,
87        builtin_table_update: Vec<BuiltinTableUpdate<&'static BuiltinTable>>,
88    ) -> Vec<BuiltinTableUpdate<CatalogItemId>> {
89        builtin_table_update
90            .into_iter()
91            .map(|builtin_table_update| self.resolve_builtin_table_update(builtin_table_update))
92            .collect()
93    }
94
95    pub fn resolve_builtin_table_update(
96        &self,
97        BuiltinTableUpdate { id, data }: BuiltinTableUpdate<&'static BuiltinTable>,
98    ) -> BuiltinTableUpdate<CatalogItemId> {
99        let id = self.resolve_builtin_table(id);
100        BuiltinTableUpdate { id, data }
101    }
102
103    pub(super) fn pack_role_auth_update(
104        &self,
105        id: RoleId,
106        diff: Diff,
107    ) -> BuiltinTableUpdate<&'static BuiltinTable> {
108        let role_auth = self.get_role_auth(&id);
109        let role = self.get_role(&id);
110        BuiltinTableUpdate::row(
111            &*MZ_ROLE_AUTH,
112            Row::pack_slice(&[
113                Datum::String(&role_auth.role_id.to_string()),
114                Datum::UInt32(role.oid),
115                match &role_auth.password_hash {
116                    Some(hash) => Datum::String(hash),
117                    None => Datum::Null,
118                },
119                Datum::TimestampTz(
120                    mz_ore::now::to_datetime(role_auth.updated_at)
121                        .try_into()
122                        .expect("must fit"),
123                ),
124            ]),
125            diff,
126        )
127    }
128
129    pub(super) fn pack_item_update(
130        &self,
131        id: CatalogItemId,
132        diff: Diff,
133    ) -> Vec<BuiltinTableUpdate<&'static BuiltinTable>> {
134        let entry = self.get_entry(&id);
135        let oid = entry.oid();
136        let conn_id = entry.item().conn_id().unwrap_or(&SYSTEM_CONN_ID);
137        let schema_id = &self
138            .get_schema(
139                &entry.name().qualifiers.database_spec,
140                &entry.name().qualifiers.schema_spec,
141                conn_id,
142            )
143            .id;
144        let name = &entry.name().item;
145        let owner_id = entry.owner_id();
146        let privileges_row = self.pack_privilege_array_row(entry.privileges());
147        let privileges = privileges_row.unpack_first();
148        let mut updates = match entry.item() {
149            CatalogItem::Index(index) => self.pack_index_update(id, index, diff),
150            CatalogItem::Source(source) => {
151                match &source.data_source {
152                    DataSourceDesc::Webhook { .. } => {
153                        vec![self.pack_webhook_source_update(id, diff)]
154                    }
155                    // Old-syntax subsource metadata (mz_postgres/mysql/sql_server_source_tables)
156                    // is now derived from create_sql by materialized views over
157                    // mz_catalog_raw, so ingestion exports need no special packing.
158                    DataSourceDesc::Ingestion { .. }
159                    | DataSourceDesc::OldSyntaxIngestion { .. }
160                    | DataSourceDesc::IngestionExport { .. }
161                    | DataSourceDesc::Introspection(_)
162                    | DataSourceDesc::Progress
163                    | DataSourceDesc::Catalog => vec![],
164                }
165            }
166            CatalogItem::MaterializedView(mview) => {
167                self.pack_materialized_view_update(id, mview, diff)
168            }
169            // mz_sinks, mz_kafka_sinks and mz_iceberg_sinks read create_sql
170            // out of mz_catalog_raw, so there is nothing to pack here.
171            CatalogItem::Sink(_) => vec![],
172            CatalogItem::Type(ty) => {
173                self.pack_type_update(id, oid, schema_id, name, owner_id, privileges, ty, diff)
174            }
175            CatalogItem::Func(func) => {
176                self.pack_func_update(id, schema_id, name, owner_id, func, diff)
177            }
178            // Tables, views, and metric sinks are exposed through materialized
179            // views derived from `mz_catalog_raw`, and logs and secrets never
180            // had builtin-table rows, so none pack a row here.
181            CatalogItem::Table(_)
182            | CatalogItem::View(_)
183            | CatalogItem::Log(_)
184            | CatalogItem::Secret(_)
185            | CatalogItem::MetricSink(_) => vec![],
186            // Connection details (mz_kafka_connections, mz_ssh_tunnel_connections,
187            // mz_aws_connections, mz_aws_privatelink_connections) are now derived
188            // from the persisted create_sql by materialized views over
189            // mz_catalog_raw, so connections need no special packing here.
190            CatalogItem::Connection(_) => vec![],
191        };
192
193        // Always report the latest for an objects columns.
194        if let Some(desc) = entry.relation_desc_latest() {
195            let defaults = match entry.item() {
196                CatalogItem::Table(Table {
197                    data_source: TableDataSource::TableWrites { defaults },
198                    ..
199                }) => Some(defaults),
200                _ => None,
201            };
202            for (i, (column_name, column_type)) in desc.iter().enumerate() {
203                let default: Option<String> = defaults.map(|d| d[i].to_ast_string_stable());
204                let default: Datum = default
205                    .as_ref()
206                    .map(|d| Datum::String(d))
207                    .unwrap_or(Datum::Null);
208                let pgtype = mz_pgrepr::Type::from(&column_type.scalar_type);
209                let (type_name, type_oid) = match &column_type.scalar_type {
210                    SqlScalarType::List {
211                        custom_id: Some(custom_id),
212                        ..
213                    }
214                    | SqlScalarType::Map {
215                        custom_id: Some(custom_id),
216                        ..
217                    }
218                    | SqlScalarType::Record {
219                        custom_id: Some(custom_id),
220                        ..
221                    } => {
222                        let entry = self.get_entry(custom_id);
223                        // NOTE(benesch): the `mz_columns.type text` field is
224                        // wrong. Types do not have a name that can be
225                        // represented as a single textual field. There can be
226                        // multiple types with the same name in different
227                        // schemas and databases. We should eventually deprecate
228                        // the `type` field in favor of a new `type_id` field
229                        // that can be joined against `mz_types`.
230                        //
231                        // For now, in the interest of pragmatism, we just use
232                        // the type's item name, and accept that there may be
233                        // ambiguity if the same type name is used in multiple
234                        // schemas. The ambiguity is mitigated by the OID, which
235                        // can be joined against `mz_types.oid` to resolve the
236                        // ambiguity.
237                        let name = &*entry.name().item;
238                        let oid = entry.oid();
239                        (name, oid)
240                    }
241                    _ => (pgtype.name(), pgtype.oid()),
242                };
243                updates.push(BuiltinTableUpdate::row(
244                    &*MZ_COLUMNS,
245                    Row::pack_slice(&[
246                        Datum::String(&id.to_string()),
247                        Datum::String(column_name),
248                        Datum::UInt64(u64::cast_from(i + 1)),
249                        Datum::from(column_type.nullable),
250                        Datum::String(type_name),
251                        default,
252                        Datum::UInt32(type_oid),
253                        Datum::Int32(pgtype.typmod()),
254                    ]),
255                    diff,
256                ));
257            }
258        }
259
260        // Use initial lcw so that we can tell apart default from non-existent windows.
261        if let Some(cw) = entry.item().initial_logical_compaction_window() {
262            updates.push(self.pack_history_retention_strategy_update(id, cw, diff));
263        }
264
265        updates
266    }
267
268    fn pack_history_retention_strategy_update(
269        &self,
270        id: CatalogItemId,
271        cw: CompactionWindow,
272        diff: Diff,
273    ) -> BuiltinTableUpdate<&'static BuiltinTable> {
274        let cw: u64 = cw.comparable_timestamp().into();
275        let cw = Jsonb::from_serde_json(serde_json::Value::Number(serde_json::Number::from(cw)))
276            .expect("must serialize");
277        BuiltinTableUpdate::row(
278            &*MZ_HISTORY_RETENTION_STRATEGIES,
279            Row::pack_slice(&[
280                Datum::String(&id.to_string()),
281                // FOR is the only strategy at the moment. We may introduce FROM or others later.
282                Datum::String("FOR"),
283                cw.into_row().into_element(),
284            ]),
285            diff,
286        )
287    }
288
289    fn pack_materialized_view_update(
290        &self,
291        id: CatalogItemId,
292        mview: &MaterializedView,
293        diff: Diff,
294    ) -> Vec<BuiltinTableUpdate<&'static BuiltinTable>> {
295        let mut updates = Vec::new();
296
297        if let Some(refresh_schedule) = &mview.refresh_schedule {
298            // This can't be `ON COMMIT`, because that is represented by a `None` instead of an
299            // empty `RefreshSchedule`.
300            assert!(!refresh_schedule.is_empty());
301            for RefreshEvery {
302                interval,
303                aligned_to,
304            } in refresh_schedule.everies.iter()
305            {
306                let aligned_to_dt = mz_ore::now::to_datetime(
307                    <&Timestamp as TryInto<u64>>::try_into(aligned_to).expect("undoes planning"),
308                );
309                updates.push(BuiltinTableUpdate::row(
310                    &*MZ_MATERIALIZED_VIEW_REFRESH_STRATEGIES,
311                    Row::pack_slice(&[
312                        Datum::String(&id.to_string()),
313                        Datum::String("every"),
314                        Datum::Interval(
315                            Interval::from_duration(interval).expect(
316                                "planning ensured that this is convertible back to Interval",
317                            ),
318                        ),
319                        Datum::TimestampTz(aligned_to_dt.try_into().expect("undoes planning")),
320                        Datum::Null,
321                    ]),
322                    diff,
323                ));
324            }
325            for at in refresh_schedule.ats.iter() {
326                let at_dt = mz_ore::now::to_datetime(
327                    <&Timestamp as TryInto<u64>>::try_into(at).expect("undoes planning"),
328                );
329                updates.push(BuiltinTableUpdate::row(
330                    &*MZ_MATERIALIZED_VIEW_REFRESH_STRATEGIES,
331                    Row::pack_slice(&[
332                        Datum::String(&id.to_string()),
333                        Datum::String("at"),
334                        Datum::Null,
335                        Datum::Null,
336                        Datum::TimestampTz(at_dt.try_into().expect("undoes planning")),
337                    ]),
338                    diff,
339                ));
340            }
341        } else {
342            updates.push(BuiltinTableUpdate::row(
343                &*MZ_MATERIALIZED_VIEW_REFRESH_STRATEGIES,
344                Row::pack_slice(&[
345                    Datum::String(&id.to_string()),
346                    Datum::String("on-commit"),
347                    Datum::Null,
348                    Datum::Null,
349                    Datum::Null,
350                ]),
351                diff,
352            ));
353        }
354
355        if let Some(target_id) = mview.replacement_target {
356            updates.push(BuiltinTableUpdate::row(
357                &*MZ_REPLACEMENTS,
358                Row::pack_slice(&[
359                    Datum::String(&id.to_string()),
360                    Datum::String(&target_id.to_string()),
361                ]),
362                diff,
363            ));
364        }
365
366        updates
367    }
368
369    fn pack_index_update(
370        &self,
371        id: CatalogItemId,
372        index: &Index,
373        diff: Diff,
374    ) -> Vec<BuiltinTableUpdate<&'static BuiltinTable>> {
375        let mut updates = vec![];
376
377        let create_stmt = mz_sql::parse::parse(&index.create_sql)
378            .unwrap_or_else(|e| {
379                panic!(
380                    "create_sql cannot be invalid: `{}` --- error: `{}`",
381                    index.create_sql, e
382                )
383            })
384            .into_element()
385            .ast;
386
387        let key_sqls = match &create_stmt {
388            Statement::CreateIndex(CreateIndexStatement { key_parts, .. }) => key_parts
389                .as_ref()
390                .expect("key_parts is filled in during planning"),
391            _ => unreachable!(),
392        };
393
394        let on_entry = self.get_entry_by_global_id(&index.on);
395        let on_desc = on_entry
396            .relation_desc()
397            .expect("can only create indexes on items with a valid description");
398        let repr_col_types: Vec<ReprColumnType> = on_desc
399            .typ()
400            .column_types
401            .iter()
402            .map(ReprColumnType::from)
403            .collect();
404        for (i, key) in index.keys.iter().enumerate() {
405            let nullable = key.typ(&repr_col_types).nullable;
406            let seq_in_index = u64::cast_from(i + 1);
407            let key_sql = key_sqls
408                .get(i)
409                .expect("missing sql information for index key")
410                .to_ast_string_simple();
411            let (field_number, expression) = match key {
412                MirScalarExpr::Column(col, _) => {
413                    (Datum::UInt64(u64::cast_from(*col + 1)), Datum::Null)
414                }
415                _ => (Datum::Null, Datum::String(&key_sql)),
416            };
417            updates.push(BuiltinTableUpdate::row(
418                &*MZ_INDEX_COLUMNS,
419                Row::pack_slice(&[
420                    Datum::String(&id.to_string()),
421                    Datum::UInt64(seq_in_index),
422                    field_number,
423                    expression,
424                    Datum::from(nullable),
425                ]),
426                diff,
427            ));
428        }
429
430        updates
431    }
432
433    fn pack_type_update(
434        &self,
435        id: CatalogItemId,
436        oid: u32,
437        schema_id: &SchemaSpecifier,
438        name: &str,
439        owner_id: &RoleId,
440        privileges: Datum,
441        typ: &Type,
442        diff: Diff,
443    ) -> Vec<BuiltinTableUpdate<&'static BuiltinTable>> {
444        let mut out = vec![];
445
446        let redacted = typ.create_sql.as_ref().map(|create_sql| {
447            mz_sql::parse::parse(create_sql)
448                .unwrap_or_else(|_| panic!("create_sql cannot be invalid: {}", create_sql))
449                .into_element()
450                .ast
451                .to_ast_string_redacted()
452        });
453
454        out.push(BuiltinTableUpdate::row(
455            &*MZ_TYPES,
456            Row::pack_slice(&[
457                Datum::String(&id.to_string()),
458                Datum::UInt32(oid),
459                Datum::String(&schema_id.to_string()),
460                Datum::String(name),
461                Datum::String(&TypeCategory::from_catalog_type(&typ.details.typ).to_string()),
462                Datum::String(&owner_id.to_string()),
463                privileges,
464                if let Some(create_sql) = &typ.create_sql {
465                    Datum::String(create_sql)
466                } else {
467                    Datum::Null
468                },
469                if let Some(redacted) = &redacted {
470                    Datum::String(redacted)
471                } else {
472                    Datum::Null
473                },
474            ]),
475            diff,
476        ));
477
478        let mut row = Row::default();
479        let mut packer = row.packer();
480
481        fn append_modifier(packer: &mut RowPacker<'_>, mods: &[i64]) {
482            if mods.is_empty() {
483                packer.push(Datum::Null);
484            } else {
485                packer.push_list(mods.iter().map(|m| Datum::Int64(*m)));
486            }
487        }
488
489        let index_id = match &typ.details.typ {
490            CatalogType::Array {
491                element_reference: element_id,
492            } => {
493                packer.push(Datum::String(&id.to_string()));
494                packer.push(Datum::String(&element_id.to_string()));
495                &MZ_ARRAY_TYPES
496            }
497            CatalogType::List {
498                element_reference: element_id,
499                element_modifiers,
500            } => {
501                packer.push(Datum::String(&id.to_string()));
502                packer.push(Datum::String(&element_id.to_string()));
503                append_modifier(&mut packer, element_modifiers);
504                &MZ_LIST_TYPES
505            }
506            CatalogType::Map {
507                key_reference: key_id,
508                value_reference: value_id,
509                key_modifiers,
510                value_modifiers,
511            } => {
512                packer.push(Datum::String(&id.to_string()));
513                packer.push(Datum::String(&key_id.to_string()));
514                packer.push(Datum::String(&value_id.to_string()));
515                append_modifier(&mut packer, key_modifiers);
516                append_modifier(&mut packer, value_modifiers);
517                &MZ_MAP_TYPES
518            }
519            CatalogType::Pseudo => {
520                packer.push(Datum::String(&id.to_string()));
521                &MZ_PSEUDO_TYPES
522            }
523            _ => {
524                packer.push(Datum::String(&id.to_string()));
525                &MZ_BASE_TYPES
526            }
527        };
528        out.push(BuiltinTableUpdate::row(index_id, row, diff));
529
530        if let Some(pg_metadata) = &typ.details.pg_metadata {
531            out.push(BuiltinTableUpdate::row(
532                &*MZ_TYPE_PG_METADATA,
533                Row::pack_slice(&[
534                    Datum::String(&id.to_string()),
535                    Datum::UInt32(pg_metadata.typinput_oid),
536                    Datum::UInt32(pg_metadata.typreceive_oid),
537                    Datum::UInt32(pg_metadata.typsend_oid),
538                ]),
539                diff,
540            ));
541        }
542
543        out
544    }
545
546    fn pack_func_update(
547        &self,
548        id: CatalogItemId,
549        schema_id: &SchemaSpecifier,
550        name: &str,
551        owner_id: &RoleId,
552        func: &Func,
553        diff: Diff,
554    ) -> Vec<BuiltinTableUpdate<&'static BuiltinTable>> {
555        let mut updates = vec![];
556        for func_impl_details in func.inner.func_impls() {
557            let arg_type_ids = func_impl_details
558                .arg_typs
559                .iter()
560                .map(|typ| self.get_system_type(typ).id().to_string())
561                .collect::<Vec<_>>();
562
563            let mut row = Row::default();
564            row.packer()
565                .try_push_array(
566                    &[ArrayDimension {
567                        lower_bound: 1,
568                        length: arg_type_ids.len(),
569                    }],
570                    arg_type_ids.iter().map(|id| Datum::String(id)),
571                )
572                .expect(
573                    "arg_type_ids is 1 dimensional, and its length is used for the array length",
574                );
575            let arg_type_ids = row.unpack_first();
576
577            updates.push(BuiltinTableUpdate::row(
578                &*MZ_FUNCTIONS,
579                Row::pack_slice(&[
580                    Datum::String(&id.to_string()),
581                    Datum::UInt32(func_impl_details.oid),
582                    Datum::String(&schema_id.to_string()),
583                    Datum::String(name),
584                    arg_type_ids,
585                    Datum::from(
586                        func_impl_details
587                            .variadic_typ
588                            .map(|typ| self.get_system_type(typ).id().to_string())
589                            .as_deref(),
590                    ),
591                    Datum::from(
592                        func_impl_details
593                            .return_typ
594                            .map(|typ| self.get_system_type(typ).id().to_string())
595                            .as_deref(),
596                    ),
597                    func_impl_details.return_is_set.into(),
598                    Datum::String(&owner_id.to_string()),
599                ]),
600                diff,
601            ));
602
603            if let mz_sql::func::Func::Aggregate(_) = func.inner {
604                updates.push(BuiltinTableUpdate::row(
605                    &*MZ_AGGREGATES,
606                    Row::pack_slice(&[
607                        Datum::UInt32(func_impl_details.oid),
608                        // TODO(database-issues#1064): Support ordered-set aggregate functions.
609                        Datum::String("n"),
610                        Datum::Int16(0),
611                    ]),
612                    diff,
613                ));
614            }
615        }
616        updates
617    }
618
619    pub fn pack_op_update(
620        &self,
621        operator: &str,
622        func_impl_details: FuncImplCatalogDetails,
623        diff: Diff,
624    ) -> BuiltinTableUpdate<&'static BuiltinTable> {
625        let arg_type_ids = func_impl_details
626            .arg_typs
627            .iter()
628            .map(|typ| self.get_system_type(typ).id().to_string())
629            .collect::<Vec<_>>();
630
631        let mut row = Row::default();
632        row.packer()
633            .try_push_array(
634                &[ArrayDimension {
635                    lower_bound: 1,
636                    length: arg_type_ids.len(),
637                }],
638                arg_type_ids.iter().map(|id| Datum::String(id)),
639            )
640            .expect("arg_type_ids is 1 dimensional, and its length is used for the array length");
641        let arg_type_ids = row.unpack_first();
642
643        BuiltinTableUpdate::row(
644            &*MZ_OPERATORS,
645            Row::pack_slice(&[
646                Datum::UInt32(func_impl_details.oid),
647                Datum::String(operator),
648                arg_type_ids,
649                Datum::from(
650                    func_impl_details
651                        .return_typ
652                        .map(|typ| self.get_system_type(typ).id().to_string())
653                        .as_deref(),
654                ),
655            ]),
656            diff,
657        )
658    }
659
660    pub fn pack_storage_usage_update(
661        &self,
662        VersionedStorageUsage::V1(event): VersionedStorageUsage,
663        diff: Diff,
664    ) -> BuiltinTableUpdate<&'static BuiltinTable> {
665        let id = &MZ_STORAGE_USAGE_BY_SHARD;
666        let row = Row::pack_slice(&[
667            Datum::UInt64(event.id),
668            Datum::from(event.shard_id.as_deref()),
669            Datum::UInt64(event.size_bytes),
670            Datum::TimestampTz(
671                mz_ore::now::to_datetime(event.collection_timestamp)
672                    .try_into()
673                    .expect("must fit"),
674            ),
675        ]);
676        BuiltinTableUpdate::row(id, row, diff)
677    }
678
679    pub fn pack_egress_ip_update(
680        &self,
681        ip: &IpNet,
682    ) -> Result<BuiltinTableUpdate<&'static BuiltinTable>, Error> {
683        let id = &MZ_EGRESS_IPS;
684        let addr = ip.network();
685        let row = Row::pack_slice(&[
686            Datum::String(&addr.to_string()),
687            Datum::Int32(ip.prefix_len().into()),
688            Datum::String(&format!("{}/{}", addr, ip.prefix_len())),
689        ]);
690        Ok(BuiltinTableUpdate::row(id, row, Diff::ONE))
691    }
692
693    pub fn pack_license_key_update(
694        &self,
695        license_key: &ValidatedLicenseKey,
696    ) -> Result<BuiltinTableUpdate<&'static BuiltinTable>, Error> {
697        let id = &MZ_LICENSE_KEYS;
698        let row = Row::pack_slice(&[
699            Datum::String(&license_key.id),
700            Datum::String(&license_key.organization),
701            Datum::String(&license_key.environment_id),
702            Datum::TimestampTz(
703                mz_ore::now::to_datetime(license_key.expiration * 1000)
704                    .try_into()
705                    .expect("must fit"),
706            ),
707            Datum::TimestampTz(
708                mz_ore::now::to_datetime(license_key.not_before * 1000)
709                    .try_into()
710                    .expect("must fit"),
711            ),
712        ]);
713        Ok(BuiltinTableUpdate::row(id, row, Diff::ONE))
714    }
715
716    pub fn pack_all_replica_size_updates(&self) -> Vec<BuiltinTableUpdate<&'static BuiltinTable>> {
717        let mut updates = Vec::new();
718        for (size, alloc) in &self.cluster_replica_sizes.0 {
719            // Just invent something when the limits are `None`, which only happens in non-prod
720            // environments (tests, process orchestrator, etc.)
721            let DiskLimit(ByteSize(disk_bytes)) =
722                (alloc.disk_limit).unwrap_or(DiskLimit::ARBITRARY);
723
724            // The disk column of mz_clusters / mz_cluster_replicas MVs needs
725            // `swap_enabled` and `disk_bytes`; expose them through a parallel
726            // internal table. Unlike the public sizes table below, we write
727            // here unconditionally — `cluster_replica_size_has_disk` previously
728            // indexed the in-memory map without checking `disabled`, so a
729            // managed cluster pinned to a disabled size still resolved its
730            // `disk` column from the real allocation. Writing disabled rows
731            // here preserves that behavior.
732            let internal_row = Row::pack_slice(&[
733                size.as_str().into(),
734                Datum::from(alloc.swap_enabled),
735                disk_bytes.into(),
736            ]);
737            updates.push(BuiltinTableUpdate::row(
738                &*MZ_CLUSTER_REPLICA_SIZE_INTERNAL,
739                internal_row,
740                Diff::ONE,
741            ));
742
743            if alloc.disabled {
744                continue;
745            }
746
747            let cpu_limit = alloc.cpu_limit.unwrap_or(CpuLimit::MAX);
748            let MemoryLimit(ByteSize(memory_bytes)) =
749                (alloc.memory_limit).unwrap_or(MemoryLimit::MAX);
750
751            let row = Row::pack_slice(&[
752                size.as_str().into(),
753                u64::cast_from(alloc.scale).into(),
754                u64::cast_from(alloc.workers).into(),
755                cpu_limit.as_nanocpus().into(),
756                memory_bytes.into(),
757                disk_bytes.into(),
758                (alloc.credits_per_hour).into(),
759            ]);
760
761            updates.push(BuiltinTableUpdate::row(
762                &*MZ_CLUSTER_REPLICA_SIZES,
763                row,
764                Diff::ONE,
765            ));
766        }
767
768        updates
769    }
770
771    pub fn pack_subscribe_update(
772        &self,
773        id: GlobalId,
774        subscribe: &ActiveSubscribe,
775        session_uuid: Uuid,
776        diff: Diff,
777    ) -> BuiltinTableUpdate<&'static BuiltinTable> {
778        let mut row = Row::default();
779        let mut packer = row.packer();
780        packer.push(Datum::String(&id.to_string()));
781        packer.push(Datum::Uuid(session_uuid));
782        packer.push(Datum::String(&subscribe.cluster_id.to_string()));
783
784        let start_dt = mz_ore::now::to_datetime(subscribe.start_time);
785        packer.push(Datum::TimestampTz(start_dt.try_into().expect("must fit")));
786
787        let depends_on: Vec<_> = subscribe
788            .depends_on
789            .iter()
790            .map(|id| id.to_string())
791            .collect();
792        packer.push_list(depends_on.iter().map(|s| Datum::String(s)));
793
794        BuiltinTableUpdate::row(&*MZ_SUBSCRIPTIONS, row, diff)
795    }
796
797    pub fn pack_session_update(
798        &self,
799        conn: &ConnMeta,
800        diff: Diff,
801    ) -> BuiltinTableUpdate<&'static BuiltinTable> {
802        let connect_dt = mz_ore::now::to_datetime(conn.connected_at());
803        BuiltinTableUpdate::row(
804            &*MZ_SESSIONS,
805            Row::pack_slice(&[
806                Datum::Uuid(conn.uuid()),
807                Datum::UInt32(conn.conn_id().unhandled()),
808                Datum::String(&conn.authenticated_role_id().to_string()),
809                Datum::from(conn.client_ip().map(|ip| ip.to_string()).as_deref()),
810                Datum::TimestampTz(connect_dt.try_into().expect("must fit")),
811            ]),
812            diff,
813        )
814    }
815
816    fn pack_privilege_array_row(&self, privileges: &PrivilegeMap) -> Row {
817        let mut row = Row::default();
818        let flat_privileges: Vec<_> = privileges.all_values_owned().collect();
819        row.packer()
820            .try_push_array(
821                &[ArrayDimension {
822                    lower_bound: 1,
823                    length: flat_privileges.len(),
824                }],
825                flat_privileges
826                    .into_iter()
827                    .map(|mz_acl_item| Datum::MzAclItem(mz_acl_item.clone())),
828            )
829            .expect("privileges is 1 dimensional, and its length is used for the array length");
830        row
831    }
832
833    pub fn pack_webhook_source_update(
834        &self,
835        item_id: CatalogItemId,
836        diff: Diff,
837    ) -> BuiltinTableUpdate<&'static BuiltinTable> {
838        let url = self
839            .try_get_webhook_url(&item_id)
840            .expect("webhook source should exist");
841        let url = url.to_string();
842        let name = &self.get_entry(&item_id).name().item;
843        let id_str = item_id.to_string();
844
845        BuiltinTableUpdate::row(
846            &*MZ_WEBHOOKS_SOURCES,
847            Row::pack_slice(&[
848                Datum::String(&id_str),
849                Datum::String(name),
850                Datum::String(&url),
851            ]),
852            diff,
853        )
854    }
855}