1use 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#[derive(clap::Parser)]
54#[clap(name = "clusterd", version = VERSION.as_str())]
55struct Args {
56 #[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 #[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 #[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 #[clap(long, env = "GRPC_HOST", value_name = "NAME")]
88 grpc_host: Option<String>,
89
90 #[clap(long, env = "STORAGE_TIMELY_CONFIG")]
93 storage_timely_config: TimelyConfig,
94 #[clap(long, env = "COMPUTE_TIMELY_CONFIG")]
96 compute_timely_config: TimelyConfig,
97 #[clap(long, env = "PROCESS")]
99 process: usize,
100
101 #[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 #[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 #[clap(long, env = "AWS_CONNECTION_ROLE_ARN")]
121 aws_connection_role_arn: Option<String>,
122
123 #[clap(flatten)]
125 secrets: SecretsReaderCliArgs,
126
127 #[clap(flatten)]
129 tracing: TracingCliArgs,
130
131 #[clap(long, env = "ENVIRONMENT_ID")]
135 environment_id: String,
136
137 #[clap(long, env = "SCRATCH_DIRECTORY", value_name = "PATH")]
139 scratch_directory: Option<PathBuf>,
140
141 #[clap(long)]
146 announce_memory_limit: Option<usize>,
147
148 #[clap(long)]
155 heap_limit: Option<usize>,
156
157 #[clap(long)]
160 is_cc: bool,
161
162 #[clap(long)]
168 worker_core_affinity: bool,
169
170 #[clap(long, env = "UNIFIED_CLUSTER")]
174 unified_cluster: bool,
175}
176
177fn 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 let _ = rustls::crypto::aws_lc_rs::default_provider().install_default();
200
201 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 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) .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 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 async move {
348 loop {
349 let (conn, remote_addr) = match listener.accept().await {
350 Ok(peer) => peer,
351 Err(error) => {
352 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 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 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 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 return future::pending().await;
493 }
494
495 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 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 future::pending().await
560}
561
562fn 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 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 assert_eq!(process_ordinal_from_hostname("7"), Some("7"));
591
592 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}