1use std::ffi::OsString;
13use std::fmt;
14use std::sync::Arc;
15use std::time::Duration;
16
17use async_trait::async_trait;
18use clap::{CommandFactory, FromArgMatches};
19use derivative::Derivative;
20use futures_core::stream::BoxStream;
21use http::header::{HeaderName, HeaderValue};
22use mz_build_info::BuildInfo;
23#[cfg(feature = "tokio-console")]
24use mz_orchestrator::ServicePort;
25use mz_orchestrator::{
26 NamespacedOrchestrator, Orchestrator, Service, ServiceAssignments, ServiceConfig, ServiceEvent,
27 ServiceProcessMetrics,
28};
29use mz_ore::cli::KeyValueArg;
30use mz_ore::metrics::MetricsRegistry;
31#[cfg(feature = "tokio-console")]
32use mz_ore::netio::SocketAddr;
33#[cfg(feature = "tokio-console")]
34use mz_ore::tracing::TokioConsoleConfig;
35use mz_ore::tracing::{
36 OpenTelemetryConfig, SentryConfig, StderrLogConfig, StderrLogFormat, TracingConfig,
37 TracingHandle,
38};
39use mz_tracing::CloneableEnvFilter;
40use opentelemetry::KeyValue;
41use opentelemetry_sdk::resource::Resource;
42use sentry_tracing::EventFilter;
43
44#[derive(Derivative, Clone, clap::Parser)]
54#[derivative(Debug)]
55pub struct TracingCliArgs {
56 #[clap(
89 long,
90 env = "STARTUP_LOG_FILTER",
91 value_name = "FILTER",
92 default_value = "info"
93 )]
94 pub startup_log_filter: CloneableEnvFilter,
95 #[clap(long, env = "LOG_FORMAT", default_value_t, value_enum)]
97 pub log_format: LogFormat,
98 #[clap(long, env = "LOG_PREFIX")]
102 pub log_prefix: Option<String>,
103 #[clap(
109 long,
110 env = "OPENTELEMETRY_MAX_BATCH_QUEUE_SIZE",
111 default_value = "2048",
112 requires = "opentelemetry_endpoint"
113 )]
114 pub opentelemetry_max_batch_queue_size: usize,
115 #[clap(
117 long,
118 env = "OPENTELEMETRY_MAX_EXPORT_BATCH_SIZE",
119 default_value = "512",
120 requires = "opentelemetry_endpoint"
121 )]
122 pub opentelemetry_max_export_batch_size: usize,
123 #[clap(
125 long,
126 env = "OPENTELEMETRY_MAX_CONCURRENT_EXPORTS",
127 default_value = "1",
128 requires = "opentelemetry_endpoint"
129 )]
130 pub opentelemetry_max_concurrent_exports: usize,
131 #[clap(
133 long,
134 env = "OPENTELEMETRY_SCHED_DELAY",
135 default_value = "5000ms",
136 requires = "opentelemetry_endpoint",
137 value_parser = humantime::parse_duration,
138 )]
139 pub opentelemetry_sched_delay: Duration,
140 #[clap(
142 long,
143 env = "OPENTELEMETRY_MAX_EXPORT_TIMEOUT",
144 default_value = "30s",
145 requires = "opentelemetry_endpoint",
146 value_parser = humantime::parse_duration,
147 )]
148 pub opentelemetry_max_export_timeout: Duration,
149 #[clap(long, env = "OPENTELEMETRY_ENDPOINT")]
155 pub opentelemetry_endpoint: Option<String>,
156 #[clap(
163 long,
164 env = "OPENTELEMETRY_HEADER",
165 requires = "opentelemetry_endpoint",
166 value_name = "NAME=VALUE",
167 use_value_delimiter = true
168 )]
169 pub opentelemetry_header: Vec<KeyValueArg<HeaderName, HeaderValue>>,
170 #[clap(
178 long,
179 env = "STARTUP_OPENTELEMETRY_FILTER",
180 requires = "opentelemetry_endpoint",
181 default_value = "info"
182 )]
183 pub startup_opentelemetry_filter: CloneableEnvFilter,
184 #[clap(
190 long,
191 env = "OPENTELEMETRY_RESOURCE",
192 value_name = "NAME=VALUE",
193 use_value_delimiter = true
194 )]
195 pub opentelemetry_resource: Vec<KeyValueArg<String, String>>,
196 #[cfg(feature = "tokio-console")]
202 #[clap(long, env = "TOKIO_CONSOLE_LISTEN_ADDR")]
203 pub tokio_console_listen_addr: Option<SocketAddr>,
204 #[cfg(feature = "tokio-console")]
208 #[clap(
209 long,
210 env = "TOKIO_CONSOLE_PUBLISH_INTERVAL",
211 requires = "tokio_console_listen_addr",
212 value_parser = humantime::parse_duration,
213 default_value = "1s",
214 )]
215 pub tokio_console_publish_interval: Duration,
216 #[cfg(feature = "tokio-console")]
220 #[clap(
221 long,
222 env = "TOKIO_CONSOLE_RETENTION",
223 requires = "tokio_console_listen_addr",
224 value_parser = humantime::parse_duration,
225 default_value = "1h",
226 )]
227 pub tokio_console_retention: Duration,
228 #[clap(long, env = "SENTRY_DSN")]
230 pub sentry_dsn: Option<String>,
231 #[clap(long, env = "SENTRY_ENVIRONMENT")]
237 pub sentry_environment: Option<String>,
238 #[clap(
243 long,
244 env = "SENTRY_TAG",
245 value_name = "NAME=VALUE",
246 use_value_delimiter = true
247 )]
248 pub sentry_tag: Vec<KeyValueArg<String, String>>,
249 #[cfg(feature = "capture")]
251 #[derivative(Debug = "ignore")]
252 #[clap(skip)]
253 pub capture: Option<tracing_capture::SharedStorage>,
254}
255
256impl Default for TracingCliArgs {
257 fn default() -> TracingCliArgs {
258 let matches = TracingCliArgs::command().get_matches_from::<_, OsString>([]);
259 TracingCliArgs::from_arg_matches(&matches)
260 .expect("no arguments produce valid TracingCliArgs")
261 }
262}
263
264impl TracingCliArgs {
265 pub async fn configure_tracing(
266 &self,
267 StaticTracingConfig {
268 service_name,
269 build_info,
270 }: StaticTracingConfig,
271 registry: MetricsRegistry,
272 ) -> Result<TracingHandle, anyhow::Error> {
273 mz_ore::tracing::configure(TracingConfig {
274 service_name,
275 stderr_log: StderrLogConfig {
276 format: match self.log_format {
277 LogFormat::Text => StderrLogFormat::Text {
278 prefix: self.log_prefix.clone(),
279 },
280 LogFormat::Json => StderrLogFormat::Json,
281 },
282 filter: self.startup_log_filter.clone().into(),
283 },
284 opentelemetry: self.opentelemetry_endpoint.clone().map(|endpoint| {
285 OpenTelemetryConfig {
286 endpoint,
287 headers: self
288 .opentelemetry_header
289 .iter()
290 .map(|header| (header.key.clone(), header.value.clone()))
291 .collect(),
292 filter: self.startup_opentelemetry_filter.clone().into(),
293 max_batch_queue_size: self.opentelemetry_max_batch_queue_size,
294 max_export_batch_size: self.opentelemetry_max_export_batch_size,
295 max_concurrent_exports: self.opentelemetry_max_concurrent_exports,
296 batch_scheduled_delay: self.opentelemetry_sched_delay,
297 max_export_timeout: self.opentelemetry_max_export_timeout,
298 resource: {
299 Resource::builder()
300 .with_attributes(
301 self.opentelemetry_resource
302 .iter()
303 .cloned()
304 .map(|kv| KeyValue::new(kv.key, kv.value)),
305 )
306 .build()
307 },
308 }
309 }),
310 #[cfg(feature = "tokio-console")]
311 tokio_console: self.tokio_console_listen_addr.clone().map(|listen_addr| {
312 TokioConsoleConfig {
313 listen_addr,
314 publish_interval: self.tokio_console_publish_interval,
315 retention: self.tokio_console_retention,
316 }
317 }),
318 sentry: self.sentry_dsn.clone().map(|dsn| SentryConfig {
319 dsn,
320 environment: self.sentry_environment.clone(),
321 tags: self
322 .opentelemetry_resource
323 .iter()
324 .cloned()
325 .chain(self.sentry_tag.iter().cloned())
326 .map(|kv| (kv.key, kv.value))
327 .collect(),
328 event_filter: mz_sentry_event_filter,
329 }),
330 build_version: build_info.version,
331 build_sha: build_info.sha,
332 registry,
333 #[cfg(feature = "capture")]
334 capture: self.capture.clone(),
335 })
336 .await
337 }
338}
339
340pub fn mz_sentry_event_filter(meta: &tracing::Metadata<'_>) -> EventFilter {
341 if meta.target() == "librdkafka" {
343 return EventFilter::Ignore;
344 }
345
346 sentry_tracing::default_event_filter(meta)
348}
349
350pub struct StaticTracingConfig {
352 pub service_name: &'static str,
354 pub build_info: BuildInfo,
356}
357
358#[derive(Debug)]
360pub struct TracingOrchestrator {
361 inner: Arc<dyn Orchestrator>,
362 tracing_args: TracingCliArgs,
363}
364
365impl TracingOrchestrator {
366 pub fn new(inner: Arc<dyn Orchestrator>, tracing_args: TracingCliArgs) -> TracingOrchestrator {
375 TracingOrchestrator {
376 inner,
377 tracing_args,
378 }
379 }
380}
381
382impl Orchestrator for TracingOrchestrator {
383 fn namespace(&self, namespace: &str) -> Arc<dyn NamespacedOrchestrator> {
384 Arc::new(NamespacedTracingOrchestrator {
385 namespace: namespace.to_string(),
386 inner: self.inner.namespace(namespace),
387 tracing_args: self.tracing_args.clone(),
388 })
389 }
390}
391
392#[derive(Debug)]
393struct NamespacedTracingOrchestrator {
394 namespace: String,
395 inner: Arc<dyn NamespacedOrchestrator>,
396 tracing_args: TracingCliArgs,
397}
398
399#[async_trait]
400impl NamespacedOrchestrator for NamespacedTracingOrchestrator {
401 async fn fetch_service_metrics(
402 &self,
403 id: &str,
404 ) -> Result<Vec<ServiceProcessMetrics>, anyhow::Error> {
405 self.inner.fetch_service_metrics(id).await
406 }
407
408 fn ensure_service(
409 &self,
410 id: &str,
411 mut service_config: ServiceConfig,
412 ) -> Result<Box<dyn Service>, anyhow::Error> {
413 let tracing_args = self.tracing_args.clone();
414 let log_prefix_arg = format!("{}-{}", self.namespace, id);
415 let args_fn = move |assigned: ServiceAssignments| {
416 #[cfg(feature = "tokio-console")]
417 let tokio_console_listen_addr = assigned.listen_addrs.get("tokio-console");
418 let mut args = (service_config.args)(assigned);
419 let TracingCliArgs {
420 startup_log_filter,
421 log_prefix,
422 log_format,
423 opentelemetry_max_batch_queue_size,
424 opentelemetry_max_export_batch_size,
425 opentelemetry_max_concurrent_exports,
426 opentelemetry_sched_delay,
427 opentelemetry_max_export_timeout,
428 opentelemetry_endpoint,
429 opentelemetry_header,
430 startup_opentelemetry_filter: _,
431 opentelemetry_resource,
432 #[cfg(feature = "tokio-console")]
433 tokio_console_listen_addr: _,
434 #[cfg(feature = "tokio-console")]
435 tokio_console_publish_interval,
436 #[cfg(feature = "tokio-console")]
437 tokio_console_retention,
438 sentry_dsn,
439 sentry_environment,
440 sentry_tag,
441 #[cfg(feature = "capture")]
442 capture: _,
443 } = &tracing_args;
444 args.push(format!("--startup-log-filter={startup_log_filter}"));
445 args.push(format!("--log-format={log_format}"));
446 if log_prefix.is_some() {
447 args.push(format!("--log-prefix={log_prefix_arg}"));
448 }
449 if let Some(endpoint) = opentelemetry_endpoint {
450 args.push(format!("--opentelemetry-endpoint={endpoint}"));
451 for kv in opentelemetry_header {
452 args.push(format!(
453 "--opentelemetry-header={}={}",
454 kv.key,
455 kv.value
456 .to_str()
457 .expect("opentelemetry-header had non-ascii value"),
458 ));
459 }
460 args.push(format!(
461 "--opentelemetry-max-batch-queue-size={opentelemetry_max_batch_queue_size}",
462 ));
463 args.push(format!(
464 "--opentelemetry-max-export-batch-size={opentelemetry_max_export_batch_size}",
465 ));
466 args.push(format!(
467 "--opentelemetry-max-concurrent-exports={opentelemetry_max_concurrent_exports}",
468 ));
469 args.push(format!(
470 "--opentelemetry-sched-delay={}ms",
471 opentelemetry_sched_delay.as_millis(),
472 ));
473 args.push(format!(
474 "--opentelemetry-max-export-timeout={}ms",
475 opentelemetry_max_export_timeout.as_millis(),
476 ));
477 }
478 #[cfg(feature = "tokio-console")]
479 if let Some(tokio_console_listen_addr) = tokio_console_listen_addr {
480 args.push(format!(
481 "--tokio-console-listen-addr={}",
482 tokio_console_listen_addr,
483 ));
484 args.push(format!(
485 "--tokio-console-publish-interval={} us",
486 tokio_console_publish_interval.as_micros(),
487 ));
488 args.push(format!(
489 "--tokio-console-retention={} us",
490 tokio_console_retention.as_micros(),
491 ));
492 }
493 if let Some(dsn) = sentry_dsn {
494 args.push(format!("--sentry-dsn={dsn}"));
495 for kv in sentry_tag {
496 args.push(format!("--sentry-tag={}={}", kv.key, kv.value));
497 }
498 }
499 if let Some(environment) = sentry_environment {
500 args.push(format!("--sentry-environment={environment}"));
501 }
502
503 if opentelemetry_endpoint.is_some() || sentry_dsn.is_some() {
504 for kv in opentelemetry_resource {
505 args.push(format!("--opentelemetry-resource={}={}", kv.key, kv.value));
506 }
507 }
508
509 args
510 };
511 service_config.args = Box::new(args_fn);
512 #[cfg(feature = "tokio-console")]
513 if self.tracing_args.tokio_console_listen_addr.is_some() {
514 service_config.ports.push(ServicePort {
515 name: "tokio-console".into(),
516 port_hint: 6669,
517 });
518 }
519 self.inner.ensure_service(id, service_config)
520 }
521
522 fn drop_service(&self, id: &str) -> Result<(), anyhow::Error> {
523 self.inner.drop_service(id)
524 }
525
526 async fn list_services(&self) -> Result<Vec<String>, anyhow::Error> {
527 self.inner.list_services().await
528 }
529
530 async fn flush(&self) -> Result<(), anyhow::Error> {
531 self.inner.flush().await
532 }
533
534 fn watch_services(&self) -> BoxStream<'static, Result<ServiceEvent, anyhow::Error>> {
535 self.inner.watch_services()
536 }
537
538 fn update_scheduling_config(
539 &self,
540 config: mz_orchestrator::scheduling_config::ServiceSchedulingConfig,
541 ) {
542 self.inner.update_scheduling_config(config)
543 }
544}
545
546#[derive(Debug, Clone, Default, clap::ValueEnum)]
548pub enum LogFormat {
549 #[default]
553 Text,
554 Json,
558}
559
560impl fmt::Display for LogFormat {
561 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
562 match self {
563 LogFormat::Text => f.write_str("text"),
564 LogFormat::Json => f.write_str("json"),
565 }
566 }
567}