Skip to main content

mz_clusterd/
lib.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
10use std::path::PathBuf;
11use std::sync::Arc;
12use std::sync::LazyLock;
13use std::time::Duration;
14
15use anyhow::Context;
16use axum::http::StatusCode;
17use axum::routing;
18use fail::FailScenario;
19use futures::future;
20use hyper_util::rt::TokioIo;
21use mz_build_info::{BuildInfo, build_info};
22use mz_cloud_resources::AwsExternalIdPrefix;
23use mz_cluster_client::client::TimelyConfig;
24use mz_compute::server::{ComputeInstanceContext, ComputeRuntimeRole};
25use mz_http_util::DynamicFilterTarget;
26use mz_orchestrator_tracing::{StaticTracingConfig, TracingCliArgs};
27use mz_ore::cli::{self, CliConfig};
28use mz_ore::error::ErrorExt;
29use mz_ore::metrics::{MetricsRegistry, register_runtime_metrics};
30use mz_ore::netio::{Listener, SocketAddr};
31use mz_ore::now::SYSTEM_TIME;
32use mz_persist_client::cache::PersistClientCache;
33use mz_persist_client::cfg::PersistConfig;
34use mz_persist_client::rpc::{GrpcPubSubClient, PersistPubSubClient, PersistPubSubClientConfig};
35use mz_service::emit_boot_diagnostics;
36use mz_service::secrets::SecretsReaderCliArgs;
37use mz_service::transport;
38use mz_service::transport::ClusterServerMetrics;
39use mz_storage::storage_state::StorageInstanceContext;
40use mz_storage_types::connections::ConnectionContext;
41use mz_txn_wal::operator::TxnsContext;
42use tokio::runtime::Handle;
43use tower::Service;
44use tracing::{Instrument, debug, error, info, info_span};
45
46mod usage_metrics;
47
48const BUILD_INFO: BuildInfo = build_info!();
49
50pub static VERSION: LazyLock<String> = LazyLock::new(|| BUILD_INFO.human_version(None));
51
52/// Independent cluster server for Materialize.
53#[derive(clap::Parser)]
54#[clap(name = "clusterd", version = VERSION.as_str())]
55struct Args {
56    // === Connection options. ===
57    /// The address on which to listen for a connection from the storage
58    /// controller.
59    #[clap(
60        long,
61        env = "STORAGE_CONTROLLER_LISTEN_ADDR",
62        value_name = "HOST:PORT",
63        default_value = "127.0.0.1:2100"
64    )]
65    storage_controller_listen_addr: SocketAddr,
66    /// The address on which to listen for a connection from the compute
67    /// controller.
68    #[clap(
69        long,
70        env = "COMPUTE_CONTROLLER_LISTEN_ADDR",
71        value_name = "HOST:PORT",
72        default_value = "127.0.0.1:2101"
73    )]
74    compute_controller_listen_addr: SocketAddr,
75    /// The address of the internal HTTP server.
76    #[clap(
77        long,
78        env = "INTERNAL_HTTP_LISTEN_ADDR",
79        value_name = "HOST:PORT",
80        default_value = "127.0.0.1:6878"
81    )]
82    internal_http_listen_addr: SocketAddr,
83    /// The FQDN of this process, for GRPC request validation.
84    ///
85    /// Not providing this value or setting it to the empty string disables host validation for
86    /// GRPC requests.
87    #[clap(long, env = "GRPC_HOST", value_name = "NAME")]
88    grpc_host: Option<String>,
89
90    // === Timely cluster options. ===
91    /// Configuration for the storage Timely cluster.
92    #[clap(long, env = "STORAGE_TIMELY_CONFIG")]
93    storage_timely_config: TimelyConfig,
94    /// Configuration for the compute Timely cluster.
95    #[clap(long, env = "COMPUTE_TIMELY_CONFIG")]
96    compute_timely_config: TimelyConfig,
97    /// The index of the process in both Timely clusters.
98    #[clap(long, env = "PROCESS")]
99    process: usize,
100
101    // === Storage options. ===
102    /// The URL for the Persist PubSub service.
103    #[clap(
104        long,
105        env = "PERSIST_PUBSUB_URL",
106        value_name = "http://HOST:PORT",
107        default_value = "http://localhost:6879"
108    )]
109    persist_pubsub_url: String,
110
111    // === Cloud options. ===
112    /// An external ID to be supplied to all AWS AssumeRole operations.
113    ///
114    /// Details: <https://docs.aws.amazon.com/IAM/latest/UserGuide/id_roles_create_for-user_externalid.html>
115    #[clap(long, env = "AWS_EXTERNAL_ID", value_name = "ID", value_parser = AwsExternalIdPrefix::new_from_cli_argument_or_environment_variable)]
116    aws_external_id_prefix: Option<AwsExternalIdPrefix>,
117
118    /// The ARN for a Materialize-controlled role to assume before assuming
119    /// a customer's requested role for an AWS connection.
120    #[clap(long, env = "AWS_CONNECTION_ROLE_ARN")]
121    aws_connection_role_arn: Option<String>,
122
123    // === Secrets reader options. ===
124    #[clap(flatten)]
125    secrets: SecretsReaderCliArgs,
126
127    // === Tracing options. ===
128    #[clap(flatten)]
129    tracing: TracingCliArgs,
130
131    // === Other options. ===
132    /// An opaque identifier for the environment in which this process is
133    /// running.
134    #[clap(long, env = "ENVIRONMENT_ID")]
135    environment_id: String,
136
137    /// A scratch directory that can be used for ephemeral storage.
138    #[clap(long, env = "SCRATCH_DIRECTORY", value_name = "PATH")]
139    scratch_directory: Option<PathBuf>,
140
141    /// Memory limit (bytes) of the cluster replica, if known.
142    ///
143    /// The limit is expected to be enforced by the orchestrator. The clusterd process only uses it
144    /// to inform configuration of backpressure mechanism.
145    #[clap(long)]
146    announce_memory_limit: Option<usize>,
147
148    /// Heap limit (bytes) of the cluster replica.
149    ///
150    /// A process heap usage is calculated as the sum of its memory and swap usage.
151    ///
152    /// In contrast to `announce_memory_limit`, this limit is enforced by the clusterd process. If
153    /// the limit is exceeded, the process terminates itself with a 167 exit code.
154    #[clap(long)]
155    heap_limit: Option<usize>,
156
157    /// Whether this size represents a modern "cc" size rather than a legacy
158    /// T-shirt size.
159    #[clap(long)]
160    is_cc: bool,
161
162    /// Set core affinity for Timely workers.
163    ///
164    /// This flag should only be set if the process is provided with exclusive access to its
165    /// supplied CPU cores. If other processes are competing over the same cores, setting core
166    /// affinity might degrade dataflow performance rather than improving it.
167    #[clap(long)]
168    worker_core_affinity: bool,
169
170    /// Host storage objects on the compute Timely cluster instead of building a separate
171    /// storage Timely cluster. The storage and compute controller protocols are served
172    /// unchanged, from the same cluster.
173    #[clap(long, env = "UNIFIED_CLUSTER")]
174    unified_cluster: bool,
175}
176
177/// The process ordinal for a StatefulSet pod, taken from the trailing
178/// `-`-delimited segment of its hostname (e.g.
179/// "mz5ncn-cluster-s1-replica-s1-gen-1-0" → "0"). This mirrors how
180/// orchestrator-kubernetes recovers the process id from pod names.
181///
182/// Returns `None` when the trailing segment is not a non-negative integer, so
183/// an unexpected hostname leaves `CLUSTERD_PROCESS` unset rather than set to a
184/// value that fails to parse as the process index.
185fn process_ordinal_from_hostname(hostname: &str) -> Option<&str> {
186    let ordinal = hostname.rsplit('-').next()?;
187    ordinal.parse::<usize>().ok().map(|_| ordinal)
188}
189
190pub fn main() {
191    mz_ore::panic::install_enhanced_handler();
192
193    // Pin the rustls crypto provider to aws-lc-rs. The LaunchDarkly SDK uses
194    // hyper-rustls, so building its client resolves the process-default rustls
195    // provider. The workspace also links rustls' `ring` feature (pulled by
196    // other hyper-rustls chains), and with both provider features enabled
197    // rustls cannot choose a default on its own and panics. The call is
198    // idempotent, so ignore the result.
199    let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
200
201    // Derive `CLUSTERD_PROCESS` (the process ordinal) from the pod hostname
202    // when running under Kubernetes and it was not set explicitly. The
203    // distroless image has no shell entrypoint to do this, so clusterd does it
204    // itself.
205    if std::env::var("KUBERNETES_SERVICE_HOST").is_ok()
206        && std::env::var("CLUSTERD_PROCESS").is_err()
207    {
208        if let Ok(hostname) = std::env::var("HOSTNAME") {
209            if let Some(ordinal) = process_ordinal_from_hostname(&hostname) {
210                // SAFETY: `set_var` is called before any threads are spawned.
211                // `install_enhanced_handler` above only registers a panic hook.
212                // That hook spawns a thread only on panic, which cannot happen
213                // before this call.
214                unsafe { std::env::set_var("CLUSTERD_PROCESS", ordinal) };
215            }
216        }
217    }
218
219    let args = cli::parse_args(CliConfig {
220        env_prefix: Some("CLUSTERD_"),
221        enable_version_flag: true,
222    });
223
224    let ncpus_useful = usize::max(1, std::cmp::min(num_cpus::get(), num_cpus::get_physical()));
225    let runtime = tokio::runtime::Builder::new_multi_thread()
226        .worker_threads(ncpus_useful)
227        .thread_stack_size(3 * 1024 * 1024) // 3 MiB
228        // The default thread name exceeds the Linux limit on thread name
229        // length, so pick something shorter. The maximum length is 16 including
230        // a \0 terminator. This gives us four decimals, which should be enough
231        // for most existing computers.
232        .thread_name_fn(|| {
233            use std::sync::atomic::{AtomicUsize, Ordering};
234            static ATOMIC_ID: AtomicUsize = AtomicUsize::new(0);
235            let id = ATOMIC_ID.fetch_add(1, Ordering::Relaxed);
236            format!("tokio:work-{}", id)
237        })
238        .enable_all()
239        .build()
240        .unwrap();
241    if let Err(err) = runtime.block_on(run(args)) {
242        panic!("clusterd: fatal: {}", err.display_with_causes());
243    }
244}
245
246async fn run(args: Args) -> Result<(), anyhow::Error> {
247    let metrics_registry = MetricsRegistry::new();
248    let tracing_handle = args
249        .tracing
250        .configure_tracing(
251            StaticTracingConfig {
252                service_name: "clusterd",
253                build_info: BUILD_INFO,
254            },
255            metrics_registry.clone(),
256        )
257        .await?;
258
259    let tracing_handle = Arc::new(tracing_handle);
260    register_runtime_metrics("main", Handle::current().metrics(), &metrics_registry);
261
262    // Keep this _after_ the mz_ore::tracing::configure call so that its panic
263    // hook runs _before_ the one that sends things to sentry.
264    mz_timely_util::panic::halt_on_timely_communication_panic();
265
266    let _failpoint_scenario = FailScenario::setup();
267
268    emit_boot_diagnostics!(&BUILD_INFO);
269
270    mz_alloc::register_metrics_into(&metrics_registry).await;
271    mz_metrics::register_metrics_into(
272        &metrics_registry,
273        mz_dyncfgs::all_dyncfgs(),
274        args.scratch_directory.clone(),
275    )
276    .await;
277
278    if let Some(heap_limit) = args.heap_limit {
279        mz_compute::memory_limiter::start_limiter(heap_limit, &metrics_registry);
280    } else {
281        info!("no heap limit announced; disabling memory limiter");
282    }
283
284    let secrets_reader = args
285        .secrets
286        .load()
287        .await
288        .context("loading secrets reader")?;
289
290    let usage_collector = Arc::new(usage_metrics::Collector {
291        disk_root: args.scratch_directory.clone(),
292    });
293
294    mz_ore::task::spawn(|| "clusterd_internal_http_server", {
295        let metrics_registry = metrics_registry.clone();
296        tracing::info!(
297            "serving internal HTTP server on {}",
298            args.internal_http_listen_addr
299        );
300        let listener = Listener::bind(args.internal_http_listen_addr).await?;
301        let mut make_service = mz_prof_http::router(&BUILD_INFO)
302            .route(
303                "/api/livez",
304                routing::get(mz_http_util::handle_liveness_check),
305            )
306            .route(
307                "/metrics",
308                routing::get(move |headers: axum::http::HeaderMap| async move {
309                    mz_http_util::handle_prometheus(&metrics_registry, headers).await
310                }),
311            )
312            .route("/api/tracing", routing::get(mz_http_util::handle_tracing))
313            .route(
314                "/api/opentelemetry/config",
315                routing::put({
316                    move |_: axum::Json<DynamicFilterTarget>| async {
317                        (
318                            StatusCode::BAD_REQUEST,
319                            "This endpoint has been replaced. \
320                                Use the `opentelemetry_filter` system variable."
321                                .to_string(),
322                        )
323                    }
324                }),
325            )
326            .route(
327                "/api/stderr/config",
328                routing::put({
329                    move |_: axum::Json<DynamicFilterTarget>| async {
330                        (
331                            StatusCode::BAD_REQUEST,
332                            "This endpoint has been replaced. \
333                                Use the `log_filter` system variable."
334                                .to_string(),
335                        )
336                    }
337                }),
338            )
339            .route(
340                "/api/usage-metrics",
341                routing::get(async move || axum::Json(usage_collector.collect())),
342            )
343            .into_make_service();
344
345        // Once https://github.com/tokio-rs/axum/pull/2479 lands, this can become just a call to
346        // `axum::serve`.
347        async move {
348            loop {
349                let (conn, remote_addr) = match listener.accept().await {
350                    Ok(peer) => peer,
351                    Err(error) => {
352                        // Match hyper's AddrIncoming error handling:
353                        // connection errors are per-connection and can be
354                        // skipped immediately; all other errors (e.g., EMFILE)
355                        // sleep to avoid a tight loop on resource exhaustion.
356                        if is_connection_error(&error) {
357                            debug!("accepted connection already errored: {error:#}");
358                        } else {
359                            error!("internal_http accept error: {error:#}");
360                            tokio::time::sleep(Duration::from_secs(1)).await;
361                        }
362                        continue;
363                    }
364                };
365
366                let tower_service = make_service.call(&conn).await.expect("infallible");
367                let hyper_service =
368                    hyper::service::service_fn(move |req| tower_service.clone().call(req));
369
370                mz_ore::task::spawn(
371                    || format!("clusterd_internal_http_server:{remote_addr}"),
372                    async move {
373                        if let Err(error) = hyper::server::conn::http1::Builder::new()
374                            .serve_connection(TokioIo::new(conn), hyper_service)
375                            .await
376                        {
377                            // This can happen when the client performs an unclean shutdown, so a
378                            // high severity isn't warranted. Might even downgrade this to DEBUG if
379                            // it turns out too noisy.
380                            info!("error serving internal_http connection: {error:#}");
381                        }
382                    },
383                );
384            }
385        }
386    });
387
388    let pubsub_caller_id = std::env::var("HOSTNAME")
389        .ok()
390        .or_else(|| args.tracing.log_prefix.clone())
391        .unwrap_or_default();
392    let mut persist_cfg =
393        PersistConfig::new(&BUILD_INFO, SYSTEM_TIME.clone(), mz_dyncfgs::all_dyncfgs());
394    persist_cfg.is_cc_active = args.is_cc;
395    persist_cfg.announce_memory_limit = args.announce_memory_limit;
396    // Start with compaction disabled, will get enabled once a cluster receives AllowWrites.
397    persist_cfg.disable_compaction();
398
399    let persist_clients = Arc::new(PersistClientCache::new(
400        persist_cfg,
401        &metrics_registry,
402        |persist_cfg, metrics| {
403            let cfg = PersistPubSubClientConfig {
404                url: args.persist_pubsub_url,
405                caller_id: pubsub_caller_id,
406                persist_cfg: persist_cfg.clone(),
407            };
408            GrpcPubSubClient::connect(cfg, metrics)
409        },
410    ));
411    let txns_ctx = TxnsContext::default();
412
413    let connection_context = ConnectionContext::from_cli_args(
414        args.environment_id,
415        &args.tracing.startup_log_filter,
416        args.aws_external_id_prefix,
417        args.aws_connection_role_arn,
418        secrets_reader,
419        None,
420    );
421
422    let grpc_host = args.grpc_host.and_then(|h| (!h.is_empty()).then_some(h));
423    let cluster_server_metrics = ClusterServerMetrics::register_with(&metrics_registry);
424
425    let mut storage_timely_config = args.storage_timely_config;
426    storage_timely_config.process = args.process;
427    let mut compute_timely_config = args.compute_timely_config;
428    compute_timely_config.process = args.process;
429
430    // We assume each storage worker has a corresponding compute worker that can process its logs.
431    assert_eq!(
432        storage_timely_config.workers, compute_timely_config.workers,
433        "storage and compute must have equal workers-per-process",
434    );
435
436    if args.unified_cluster {
437        info!("running with a unified timely cluster");
438
439        let (compute_client_builder, storage_client_builder) = mz_compute::server::serve_unified(
440            compute_timely_config,
441            ComputeRuntimeRole::Solo,
442            &metrics_registry,
443            persist_clients,
444            txns_ctx,
445            tracing_handle,
446            ComputeInstanceContext {
447                scratch_directory: args.scratch_directory.clone(),
448                worker_core_affinity: args.worker_core_affinity,
449                connection_context: connection_context.clone(),
450            },
451            SYSTEM_TIME.clone(),
452            connection_context,
453            StorageInstanceContext::new(args.scratch_directory, args.announce_memory_limit),
454        )
455        .await?;
456
457        info!(
458            "listening for storage controller connections on {}",
459            args.storage_controller_listen_addr
460        );
461        mz_ore::task::spawn(
462            || "storage_server",
463            transport::serve(
464                args.storage_controller_listen_addr,
465                BUILD_INFO.semver_version(),
466                grpc_host.clone(),
467                Duration::MAX,
468                storage_client_builder,
469                cluster_server_metrics.for_server("storage"),
470            )
471            .instrument(info_span!("ctp", name = "storage")),
472        );
473
474        info!(
475            "listening for compute controller connections on {}",
476            args.compute_controller_listen_addr
477        );
478        mz_ore::task::spawn(
479            || "compute_server",
480            transport::serve(
481                args.compute_controller_listen_addr,
482                BUILD_INFO.semver_version(),
483                grpc_host,
484                Duration::MAX,
485                compute_client_builder,
486                cluster_server_metrics.for_server("compute"),
487            )
488            .instrument(info_span!("ctp", name = "compute")),
489        );
490
491        // Block forever.
492        return future::pending().await;
493    }
494
495    // Start storage server.
496    let storage_client_builder = mz_storage::serve(
497        storage_timely_config,
498        &metrics_registry,
499        Arc::clone(&persist_clients),
500        txns_ctx.clone(),
501        Arc::clone(&tracing_handle),
502        SYSTEM_TIME.clone(),
503        connection_context.clone(),
504        StorageInstanceContext::new(args.scratch_directory.clone(), args.announce_memory_limit),
505    )
506    .await?;
507    info!(
508        "listening for storage controller connections on {}",
509        args.storage_controller_listen_addr
510    );
511    mz_ore::task::spawn(
512        || "storage_server",
513        transport::serve(
514            args.storage_controller_listen_addr,
515            BUILD_INFO.semver_version(),
516            grpc_host.clone(),
517            Duration::MAX,
518            storage_client_builder,
519            cluster_server_metrics.for_server("storage"),
520        )
521        .instrument(info_span!("ctp", name = "storage")),
522    );
523
524    // Start compute server.
525    let compute_client_builder = mz_compute::server::serve(
526        compute_timely_config,
527        ComputeRuntimeRole::Solo,
528        &metrics_registry,
529        persist_clients,
530        txns_ctx,
531        tracing_handle,
532        ComputeInstanceContext {
533            scratch_directory: args.scratch_directory,
534            worker_core_affinity: args.worker_core_affinity,
535            connection_context,
536        },
537    )
538    .await?;
539    info!(
540        "listening for compute controller connections on {}",
541        args.compute_controller_listen_addr
542    );
543    mz_ore::task::spawn(
544        || "compute_server",
545        transport::serve(
546            args.compute_controller_listen_addr,
547            BUILD_INFO.semver_version(),
548            grpc_host.clone(),
549            Duration::MAX,
550            compute_client_builder,
551            cluster_server_metrics.for_server("compute"),
552        )
553        .instrument(info_span!("ctp", name = "compute")),
554    );
555
556    // TODO: retire this two-cluster topology once the unified cluster has production mileage.
557
558    // Block forever.
559    future::pending().await
560}
561
562/// Per-connection errors from `accept()` that can be skipped immediately.
563/// All other errors (e.g., EMFILE/ENFILE resource exhaustion) warrant a sleep
564/// before retrying. Mirrors hyper's `AddrIncoming` classification.
565fn is_connection_error(e: &std::io::Error) -> bool {
566    matches!(
567        e.kind(),
568        std::io::ErrorKind::ConnectionRefused
569            | std::io::ErrorKind::ConnectionAborted
570            | std::io::ErrorKind::ConnectionReset
571    )
572}
573
574#[cfg(test)]
575mod tests {
576    use super::*;
577
578    #[mz_ore::test]
579    fn test_process_ordinal_from_hostname() {
580        // A StatefulSet pod name ends in the process ordinal.
581        assert_eq!(
582            process_ordinal_from_hostname("mz5ncn-cluster-s1-replica-s1-gen-1-0"),
583            Some("0")
584        );
585        assert_eq!(
586            process_ordinal_from_hostname("mz5ncn-cluster-s1-replica-s1-gen-1-11"),
587            Some("11")
588        );
589        // A bare numeric hostname is its own ordinal.
590        assert_eq!(process_ordinal_from_hostname("7"), Some("7"));
591
592        // A trailing segment that is not a non-negative integer yields `None`,
593        // so `CLUSTERD_PROCESS` stays unset rather than being set to a value
594        // that fails to parse as the process index.
595        assert_eq!(process_ordinal_from_hostname("clusterd"), None);
596        assert_eq!(process_ordinal_from_hostname("replica-abc"), None);
597        assert_eq!(process_ordinal_from_hostname("replica-"), None);
598        assert_eq!(process_ordinal_from_hostname(""), None);
599    }
600}