1use mz_dyncfg::{Config, ConfigSet, ParameterScope};
14use std::time::Duration;
15
16pub const CLUSTER_SHUTDOWN_GRACE_PERIOD: Config<Duration> = Config::new(
20 "storage_cluster_shutdown_grace_period",
21 Duration::from_secs(10 * 60),
22 "When dataflows observe an invariant violation it is either due to a bug or due to \
23 the cluster being shut down. This configuration defines the amount of time to \
24 wait before panicking the process, which will register the invariant violation.",
25 ParameterScope::Replica,
26);
27
28pub const DELAY_SOURCES_PAST_REHYDRATION: Config<bool> = Config::new(
33 "storage_dataflow_delay_sources_past_rehydration",
34 true,
36 "Whether or not to delay sources producing values in some scenarios \
37 (namely, upsert) till after rehydration is finished",
38 ParameterScope::Environment,
39);
40
41pub const SUSPENDABLE_SOURCES: Config<bool> = Config::new(
44 "storage_dataflow_suspendable_sources",
45 true,
46 "Whether storage dataflows should suspend execution while downstream operators are still \
47 processing data.",
48 ParameterScope::Environment,
49);
50
51pub const STORAGE_DOWNGRADE_SINCE_DURING_FINALIZATION: Config<bool> = Config::new(
56 "storage_downgrade_since_during_finalization",
57 true,
59 "When enabled, force-downgrade the controller's since handle on the shard\
60 during shard finalization",
61 ParameterScope::Environment,
62);
63
64pub const REPLICA_METRICS_HISTORY_RETENTION_INTERVAL: Config<Duration> = Config::new(
66 "replica_metrics_history_retention_interval",
67 Duration::from_secs(60 * 60 * 24 * 30), "The interval of time to keep when truncating the replica metrics history.",
69 ParameterScope::Environment,
70);
71
72pub const WALLCLOCK_LAG_HISTORY_RETENTION_INTERVAL: Config<Duration> = Config::new(
74 "wallclock_lag_history_retention_interval",
75 Duration::from_secs(60 * 60 * 24 * 30), "The interval of time to keep when truncating the wallclock lag history.",
77 ParameterScope::Environment,
78);
79
80pub const WALLCLOCK_GLOBAL_LAG_HISTOGRAM_RETENTION_INTERVAL: Config<Duration> = Config::new(
82 "wallclock_global_lag_histogram_retention_interval",
83 Duration::from_secs(60 * 60 * 24 * 30), "The interval of time to keep when truncating the wallclock lag histogram.",
85 ParameterScope::Environment,
86);
87
88pub const KAFKA_CLIENT_ID_ENRICHMENT_RULES: Config<fn() -> serde_json::Value> = Config::new(
101 "kafka_client_id_enrichment_rules",
102 || serde_json::json!([]),
103 "Rules for enriching the `client.id` property of Kafka clients with additional data.",
104 ParameterScope::Environment,
105);
106
107pub const KAFKA_POLL_MAX_WAIT: Config<Duration> = Config::new(
110 "kafka_poll_max_wait",
111 Duration::from_secs(1),
112 "The maximum time we will wait before re-polling rdkafka to see if new partitions/data are \
113 available.",
114 ParameterScope::Replica,
115);
116
117pub const KAFKA_LOW_WATERMARK_CHECK: Config<bool> = Config::new(
123 "kafka_low_watermark_check",
124 true,
125 "Whether to check the low watermark for Kafka sources and error if the start \
126 offset/resume upper has been compacted away.",
127 ParameterScope::Environment,
128);
129
130pub const KAFKA_DEFAULT_AWS_PRIVATELINK_ENDPOINT_IDENTIFICATION_ALGORITHM: Config<&'static str> =
131 Config::new(
132 "kafka_default_aws_privatelink_endpoint_identification_algorithm",
133 "none",
135 "The value we set for the 'ssl.endpoint.identification.algorithm' option in the Kafka \
136 Connection config. default: 'none'",
137 ParameterScope::Environment,
138 );
139
140pub const KAFKA_BUFFERED_EVENT_RESIZE_THRESHOLD_ELEMENTS: Config<usize> = Config::new(
141 "kafka_buffered_event_resize_threshold_elements",
142 1000,
143 "In the Kafka sink operator we might need to buffer messages before emitting them. As a \
144 performance optimization we reuse the buffer allocations, but shrink it to retain at \
145 most this number of elements.",
146 ParameterScope::Replica,
147);
148
149pub const KAFKA_RETRY_BACKOFF: Config<Duration> = Config::new(
152 "kafka_retry_backoff",
153 Duration::from_millis(100),
154 "Sets retry.backoff.ms in librdkafka for sources and sinks.",
155 ParameterScope::Replica,
156);
157
158pub const KAFKA_RETRY_BACKOFF_MAX: Config<Duration> = Config::new(
161 "kafka_retry_backoff_max",
162 Duration::from_secs(1),
163 "Sets retry.backoff.max.ms in librdkafka for sources and sinks.",
164 ParameterScope::Replica,
165);
166
167pub const KAFKA_RECONNECT_BACKOFF: Config<Duration> = Config::new(
170 "kafka_reconnect_backoff",
171 Duration::from_millis(100),
172 "Sets reconnect.backoff.ms in librdkafka for sources and sinks.",
173 ParameterScope::Replica,
174);
175
176pub const KAFKA_RECONNECT_BACKOFF_MAX: Config<Duration> = Config::new(
181 "kafka_reconnect_backoff_max",
182 Duration::from_secs(30),
183 "Sets reconnect.backoff.max.ms in librdkafka for sources and sinks.",
184 ParameterScope::Replica,
185);
186
187pub const KAFKA_SINK_MESSAGE_MAX_BYTES: Config<usize> = Config::new(
193 "kafka_sink_message_max_bytes",
194 1_000_000,
195 "Sets message.max.bytes in librdkafka for Kafka sink producers.",
196 ParameterScope::Environment,
197);
198
199pub const KAFKA_SINK_BATCH_SIZE: Config<usize> = Config::new(
205 "kafka_sink_batch_size",
206 1_000_000,
207 "Sets batch.size in librdkafka for Kafka sink producers.",
208 ParameterScope::Environment,
209);
210
211pub const KAFKA_SINK_BATCH_NUM_MESSAGES: Config<usize> = Config::new(
216 "kafka_sink_batch_num_messages",
217 10_000,
218 "Sets batch.num.messages in librdkafka for Kafka sink producers.",
219 ParameterScope::Environment,
220);
221
222pub const MYSQL_REPLICATION_HEARTBEAT_INTERVAL: Config<Duration> = Config::new(
226 "mysql_replication_heartbeat_interval",
227 Duration::from_secs(30),
228 "Replication heartbeat interval requested from the MySQL server.",
229 ParameterScope::Replica,
230);
231
232pub static MYSQL_SOURCE_SNAPSHOT_PARALLELISM: Config<bool> = Config::new(
236 "mysql_source_snapshot_parallelism",
237 false,
238 "Whether to split MySQL snapshot reads across workers by primary-key ranges.",
239 ParameterScope::Replica,
240);
241
242pub static MYSQL_SOURCE_SNAPSHOT_PARTITION_MIN_ROWS: Config<usize> = Config::new(
244 "mysql_source_snapshot_partition_min_rows",
245 50_000,
246 "Minimum estimated rows the MySQL snapshot partitioner attempts to split.",
247 ParameterScope::Replica,
248);
249
250pub static MYSQL_SOURCE_SNAPSHOT_PARTITION_PROBED_PREFIXES_PER_BILLION_ROWS: Config<usize> =
254 Config::new(
255 "mysql_source_snapshot_partition_probed_prefixes_per_billion_rows",
256 1_000,
257 "Cap on MySQL snapshot PK-prefix partitioning probed prefixes per table, per billion \
258 estimated rows; when exhausted, splitting stops early with coarser partition boundaries.",
259 ParameterScope::Replica,
260 );
261
262pub static MYSQL_SOURCE_SNAPSHOT_EXACT_COUNT_MAX_ROWS: Config<usize> = Config::new(
265 "mysql_source_snapshot_exact_count_max_rows",
266 1_000_000,
267 "Maximum estimated table size for which MySQL snapshots compute an exact COUNT(*) \
268 for the size gauge; larger tables report the information_schema estimate.",
269 ParameterScope::Replica,
270);
271
272pub const PG_FETCH_SLOT_RESUME_LSN_INTERVAL: Config<Duration> = Config::new(
276 "postgres_fetch_slot_resume_lsn_interval",
277 Duration::from_millis(500),
278 "Interval to poll `confirmed_flush_lsn` to get a resumption lsn.",
279 ParameterScope::Replica,
280);
281
282pub const PG_SCHEMA_VALIDATION_INTERVAL: Config<Duration> = Config::new(
284 "pg_schema_validation_interval",
285 Duration::from_secs(15),
286 "Interval to re-validate the schemas of ingested tables.",
287 ParameterScope::Environment,
288);
289
290pub static PG_SOURCE_VALIDATE_TIMELINE: Config<bool> = Config::new(
299 "pg_source_validate_timeline",
300 true,
301 "Whether to treat a timeline switch as a definite error",
302 ParameterScope::Environment,
303);
304
305pub static SQL_SERVER_SOURCE_VALIDATE_RESTORE_HISTORY: Config<bool> = Config::new(
314 "sql_server_source_validate_restore_history",
315 true,
316 "Whether to treat a restore history change as a definite error",
317 ParameterScope::Environment,
318);
319
320pub const AWS_PREFETCH_STS_CONNECT_TIMEOUT: Config<Duration> = Config::new(
331 "aws_prefetch_sts_connect_timeout",
332 Duration::from_millis(3100),
333 "Connect timeout for the AWS AssumeRole credentials prefetcher's STS calls.",
334 ParameterScope::Replica,
335);
336
337pub const ENFORCE_EXTERNAL_ADDRESSES: Config<bool> = Config::new(
346 "storage_enforce_external_addresses",
347 false,
348 "Whether or not to enforce that external connection addresses are global \
349 (not private or local) when resolving them",
350 ParameterScope::Environment,
351);
352
353pub const STORAGE_UPSERT_PREVENT_SNAPSHOT_BUFFERING: Config<bool> = Config::new(
370 "storage_upsert_prevent_snapshot_buffering",
371 true,
372 "Prevent snapshot buffering in upsert.",
373 ParameterScope::Replica,
374);
375
376pub const STORAGE_ROCKSDB_USE_MERGE_OPERATOR: Config<bool> = Config::new(
378 "storage_rocksdb_use_merge_operator",
379 true,
380 "Use the native rocksdb merge operator where possible.",
381 ParameterScope::Environment,
382);
383
384pub const STORAGE_UPSERT_MAX_SNAPSHOT_BATCH_BUFFERING: Config<Option<usize>> = Config::new(
389 "storage_upsert_max_snapshot_batch_buffering",
390 None,
391 "Limit snapshot buffering in upsert.",
392 ParameterScope::Replica,
393);
394
395pub const ENABLE_UPSERT_PAGED_SPILL: Config<bool> = Config::new(
406 "enable_upsert_paged_spill",
407 false,
408 "Allow the upsert-v2 source stash to spill chunks to the shared buffer pool, gated \
409 independently of the compute `enable_column_paged_batcher_spill`.",
410 ParameterScope::Replica,
411);
412
413pub const STORAGE_ROCKSDB_CLEANUP_TRIES: Config<usize> = Config::new(
417 "storage_rocksdb_cleanup_tries",
418 5,
419 "How many times to try to cleanup old RocksDB DB's on disk before giving up.",
420 ParameterScope::Replica,
421);
422
423pub const STORAGE_SUSPEND_AND_RESTART_DELAY: Config<Duration> = Config::new(
425 "storage_suspend_and_restart_delay",
426 Duration::from_secs(5),
427 "Delay interval when reconnecting to a source / sink after halt.",
428 ParameterScope::Replica,
429);
430
431pub const STORAGE_USE_CONTINUAL_FEEDBACK_UPSERT: Config<bool> = Config::new(
433 "storage_use_continual_feedback_upsert",
434 true,
435 "Whether to use the new continual feedback upsert operator.",
436 ParameterScope::Environment,
437);
438
439pub const ENABLE_UPSERT_V2: Config<bool> = Config::new(
441 "enable_upsert_v2",
442 false,
443 "Whether to use the v2 upsert operator.",
444 ParameterScope::Environment,
445);
446
447pub const STORAGE_SERVER_MAINTENANCE_INTERVAL: Config<Duration> = Config::new(
449 "storage_server_maintenance_interval",
450 Duration::from_millis(10),
451 "The interval at which the storage server performs maintenance tasks. Zero enables maintenance on every iteration.",
452 ParameterScope::Replica,
453);
454
455pub const SINK_PROGRESS_SEARCH: Config<bool> = Config::new(
457 "storage_sink_progress_search",
458 true,
459 "If set, iteratively search the progress topic for a progress record with increasing lookback.",
460 ParameterScope::Environment,
461);
462
463pub const SINK_ENSURE_TOPIC_CONFIG: Config<&'static str> = Config::new(
465 "storage_sink_ensure_topic_config",
466 "skip",
467 "If `skip`, don't check the config of existing topics; if `check`, fetch the config and \
468 warn if it does not match the expected configs; if `alter`, attempt to change the upstream to \
469 match the expected configs.",
470 ParameterScope::Environment,
471);
472
473pub const ORE_OVERFLOWING_BEHAVIOR: Config<&'static str> = Config::new(
475 "ore_overflowing_behavior",
476 "soft_panic",
477 "Overflow behavior for Overflowing types. One of 'ignore', 'panic', 'soft_panic'.",
478 ParameterScope::Environment,
479);
480
481pub const STATISTICS_RETENTION_DURATION: Config<Duration> = Config::new(
487 "storage_statistics_retention_duration",
488 Duration::from_secs(86_400), "The time after which we delete per replica statistics (for sources and sinks) after there have been no updates.",
490 ParameterScope::Environment,
491);
492
493pub fn all_dyncfgs(configs: ConfigSet) -> ConfigSet {
495 configs
496 .add(&AWS_PREFETCH_STS_CONNECT_TIMEOUT)
497 .add(&CLUSTER_SHUTDOWN_GRACE_PERIOD)
498 .add(&DELAY_SOURCES_PAST_REHYDRATION)
499 .add(&ENFORCE_EXTERNAL_ADDRESSES)
500 .add(&KAFKA_BUFFERED_EVENT_RESIZE_THRESHOLD_ELEMENTS)
501 .add(&KAFKA_CLIENT_ID_ENRICHMENT_RULES)
502 .add(&KAFKA_DEFAULT_AWS_PRIVATELINK_ENDPOINT_IDENTIFICATION_ALGORITHM)
503 .add(&KAFKA_LOW_WATERMARK_CHECK)
504 .add(&KAFKA_POLL_MAX_WAIT)
505 .add(&KAFKA_RETRY_BACKOFF)
506 .add(&KAFKA_RETRY_BACKOFF_MAX)
507 .add(&KAFKA_RECONNECT_BACKOFF)
508 .add(&KAFKA_RECONNECT_BACKOFF_MAX)
509 .add(&KAFKA_SINK_MESSAGE_MAX_BYTES)
510 .add(&KAFKA_SINK_BATCH_SIZE)
511 .add(&KAFKA_SINK_BATCH_NUM_MESSAGES)
512 .add(&MYSQL_REPLICATION_HEARTBEAT_INTERVAL)
513 .add(&MYSQL_SOURCE_SNAPSHOT_EXACT_COUNT_MAX_ROWS)
514 .add(&MYSQL_SOURCE_SNAPSHOT_PARALLELISM)
515 .add(&MYSQL_SOURCE_SNAPSHOT_PARTITION_MIN_ROWS)
516 .add(&MYSQL_SOURCE_SNAPSHOT_PARTITION_PROBED_PREFIXES_PER_BILLION_ROWS)
517 .add(&ORE_OVERFLOWING_BEHAVIOR)
518 .add(&PG_FETCH_SLOT_RESUME_LSN_INTERVAL)
519 .add(&PG_SCHEMA_VALIDATION_INTERVAL)
520 .add(&PG_SOURCE_VALIDATE_TIMELINE)
521 .add(&REPLICA_METRICS_HISTORY_RETENTION_INTERVAL)
522 .add(&SINK_ENSURE_TOPIC_CONFIG)
523 .add(&SINK_PROGRESS_SEARCH)
524 .add(&SQL_SERVER_SOURCE_VALIDATE_RESTORE_HISTORY)
525 .add(&STORAGE_DOWNGRADE_SINCE_DURING_FINALIZATION)
526 .add(&STORAGE_ROCKSDB_CLEANUP_TRIES)
527 .add(&STORAGE_ROCKSDB_USE_MERGE_OPERATOR)
528 .add(&STORAGE_SERVER_MAINTENANCE_INTERVAL)
529 .add(&STORAGE_SUSPEND_AND_RESTART_DELAY)
530 .add(&STORAGE_UPSERT_MAX_SNAPSHOT_BATCH_BUFFERING)
531 .add(&STORAGE_UPSERT_PREVENT_SNAPSHOT_BUFFERING)
532 .add(&STORAGE_USE_CONTINUAL_FEEDBACK_UPSERT)
533 .add(&ENABLE_UPSERT_V2)
534 .add(&SUSPENDABLE_SOURCES)
535 .add(&ENABLE_UPSERT_PAGED_SPILL)
536 .add(&WALLCLOCK_GLOBAL_LAG_HISTOGRAM_RETENTION_INTERVAL)
537 .add(&WALLCLOCK_LAG_HISTORY_RETENTION_INTERVAL)
538 .add(&crate::sources::sql_server::CDC_CLEANUP_CHANGE_TABLE)
539 .add(&crate::sources::sql_server::CDC_CLEANUP_CHANGE_TABLE_MAX_DELETES)
540 .add(&crate::sources::sql_server::MAX_LSN_WAIT)
541 .add(&crate::sources::sql_server::SNAPSHOT_PROGRESS_REPORT_INTERVAL)
542 .add(&STATISTICS_RETENTION_DURATION)
543}