1mod 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
53use crate::active_compute_sink::ActiveSubscribe;
55use crate::catalog::CatalogState;
56use crate::coord::ConnMeta;
57
58#[derive(Debug, Clone)]
60pub struct BuiltinTableUpdate<T = CatalogItemId> {
61 pub id: T,
63 pub data: TableData,
65}
66
67impl<T> BuiltinTableUpdate<T> {
68 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 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 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 CatalogItem::Table(_)
182 | CatalogItem::View(_)
183 | CatalogItem::Log(_)
184 | CatalogItem::Secret(_)
185 | CatalogItem::MetricSink(_) => vec![],
186 CatalogItem::Connection(_) => vec![],
191 };
192
193 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 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 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 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 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 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 let DiskLimit(ByteSize(disk_bytes)) =
722 (alloc.disk_limit).unwrap_or(DiskLimit::ARBITRARY);
723
724 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}