|
5 | 5 |
|
6 | 6 | pub mod processors; |
7 | 7 |
|
| 8 | +use crate::error::Error; |
8 | 9 | use schemars::JsonSchema; |
9 | 10 | use serde::{Deserialize, Serialize}; |
10 | 11 |
|
11 | 12 | /// Internal logs configuration. |
12 | | -#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, Default)] |
| 13 | +#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)] |
13 | 14 | pub struct LogsConfig { |
14 | 15 | /// The log level for internal engine logs. |
15 | 16 | #[serde(default)] |
16 | 17 | pub level: LogLevel, |
17 | 18 |
|
18 | | - /// The list of log processors to configure. |
| 19 | + /// Logging provider configuration. |
| 20 | + #[serde(default = "default_providers")] |
| 21 | + pub providers: LoggingProviders, |
| 22 | + |
| 23 | + /// OpenTelemetry SDK is configured via processors. |
19 | 24 | #[serde(default)] |
20 | 25 | pub processors: Vec<processors::LogProcessorConfig>, |
21 | 26 | } |
22 | 27 |
|
23 | | -/// Log level for internal engine logs. |
24 | | -/// |
25 | | -/// TODO: Change default to `Info` once per-thread subscriber is implemented |
26 | | -/// to avoid contention from the global tracing subscriber. |
27 | | -#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, Default, PartialEq)] |
| 28 | +/// Log level for dataflow engine logs. |
| 29 | +#[derive(Debug, Clone, Copy, Serialize, Deserialize, JsonSchema, Default, PartialEq)] |
28 | 30 | #[serde(rename_all = "lowercase")] |
29 | 31 | pub enum LogLevel { |
30 | 32 | /// Logging is completely disabled. |
31 | | - #[default] |
32 | 33 | Off, |
33 | 34 | /// Debug level logging. |
34 | 35 | Debug, |
35 | 36 | /// Info level logging. |
| 37 | + #[default] |
36 | 38 | Info, |
37 | 39 | /// Warn level logging. |
38 | 40 | Warn, |
39 | 41 | /// Error level logging. |
40 | 42 | Error, |
41 | 43 | } |
42 | 44 |
|
| 45 | +/// Logging providers for different execution contexts. |
| 46 | +#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema)] |
| 47 | +pub struct LoggingProviders { |
| 48 | + /// Provider mode for non-engine threads. This defines the global Tokio |
| 49 | + /// `tracing` subscriber. Default is ConsoleAsync. |
| 50 | + #[serde(default = "default_global_provider")] |
| 51 | + pub global: ProviderMode, |
| 52 | + |
| 53 | + /// Provider mod for engine/pipeline threads. This defines how the |
| 54 | + /// engine thread / core sets the Tokio `tracing` |
| 55 | + /// subscriber. Default is ConsoleAsync. Internal logs will be flushed |
| 56 | + /// by either the Internal Telemetry Receiver or the main pipeline |
| 57 | + /// controller. |
| 58 | + #[serde(default = "default_engine_provider")] |
| 59 | + pub engine: ProviderMode, |
| 60 | + |
| 61 | + /// Provider mode for nodes downstream of Internal Telemetry receiver. |
| 62 | + /// This defaults to Noop to avoid internal feedback. |
| 63 | + #[serde(default = "default_internal_provider")] |
| 64 | + pub internal: ProviderMode, |
| 65 | +} |
| 66 | + |
| 67 | +/// Logs producer: how log events are captured and routed. |
| 68 | +#[derive(Debug, Clone, Copy, Serialize, Deserialize, JsonSchema, PartialEq)] |
| 69 | +#[serde(rename_all = "lowercase")] |
| 70 | +pub enum ProviderMode { |
| 71 | + /// Log events are silently ignored. |
| 72 | + Noop, |
| 73 | + |
| 74 | + /// Delivery using the internal telemetry system. |
| 75 | + ITS, |
| 76 | + |
| 77 | + /// Use OTel-Rust as the provider. |
| 78 | + OpenTelemetry, |
| 79 | + |
| 80 | + /// Asynchronous console logging. The caller writes to a channel |
| 81 | + /// the same as ITS delivery, but bypasses the internal pipeline |
| 82 | + /// with console logging. |
| 83 | + #[serde(rename = "console_async")] |
| 84 | + ConsoleAsync, |
| 85 | + |
| 86 | + /// Synchronous console logging. Note! This can block the |
| 87 | + /// producing thread. The caller writes directly to the console. |
| 88 | + #[serde(rename = "console_direct")] |
| 89 | + ConsoleDirect, |
| 90 | +} |
| 91 | + |
| 92 | +impl ProviderMode { |
| 93 | + /// Returns true if this requires a LogsReporter channel for |
| 94 | + /// asynchronous logging. |
| 95 | + #[must_use] |
| 96 | + pub const fn needs_reporter(&self) -> bool { |
| 97 | + matches!(self, Self::ITS | Self::ConsoleAsync) |
| 98 | + } |
| 99 | +} |
| 100 | + |
| 101 | +const fn default_global_provider() -> ProviderMode { |
| 102 | + ProviderMode::ConsoleAsync |
| 103 | +} |
| 104 | + |
| 105 | +const fn default_engine_provider() -> ProviderMode { |
| 106 | + ProviderMode::ConsoleAsync |
| 107 | +} |
| 108 | + |
| 109 | +const fn default_internal_provider() -> ProviderMode { |
| 110 | + ProviderMode::Noop |
| 111 | +} |
| 112 | + |
| 113 | +fn default_providers() -> LoggingProviders { |
| 114 | + LoggingProviders { |
| 115 | + global: default_global_provider(), |
| 116 | + engine: default_engine_provider(), |
| 117 | + internal: default_internal_provider(), |
| 118 | + } |
| 119 | +} |
| 120 | + |
| 121 | +impl Default for LogsConfig { |
| 122 | + fn default() -> Self { |
| 123 | + Self { |
| 124 | + level: LogLevel::default(), |
| 125 | + providers: default_providers(), |
| 126 | + processors: Vec::new(), |
| 127 | + } |
| 128 | + } |
| 129 | +} |
| 130 | + |
| 131 | +impl LogsConfig { |
| 132 | + /// Validate the logs configuration. |
| 133 | + /// |
| 134 | + /// Returns an error if: |
| 135 | + /// - `internal` is configured to use ITS, ConsoleAsync (needs_reporter()) |
| 136 | + /// - `engine` is `OpenTelemetry` but `global` is not |
| 137 | + /// (current implementation restriction). |
| 138 | + pub fn validate(&self) -> Result<(), Error> { |
| 139 | + if self.providers.internal.needs_reporter() { |
| 140 | + return Err(Error::InvalidUserConfig { |
| 141 | + error: format!( |
| 142 | + "internal provider is invalid: {:?}", |
| 143 | + self.providers.internal |
| 144 | + ), |
| 145 | + }); |
| 146 | + } |
| 147 | + // Current implementation restriction: engine OpenTelemetry requires global OpenTelemetry. |
| 148 | + // The SDK logger provider is only created when the global provider is OpenTelemetry. |
| 149 | + // This could be lifted in the future by creating the logger provider independently. |
| 150 | + if self.providers.engine == ProviderMode::OpenTelemetry |
| 151 | + && self.providers.global != ProviderMode::OpenTelemetry |
| 152 | + { |
| 153 | + return Err(Error::InvalidUserConfig { |
| 154 | + error: "engine provider 'opentelemetry' requires global provider to also be \ |
| 155 | + 'opentelemetry' (current implementation restriction)" |
| 156 | + .into(), |
| 157 | + }); |
| 158 | + } |
| 159 | + |
| 160 | + Ok(()) |
| 161 | + } |
| 162 | +} |
| 163 | + |
43 | 164 | #[cfg(test)] |
44 | 165 | mod tests { |
45 | 166 | use super::*; |
46 | 167 |
|
| 168 | + /// Helper to parse YAML into LogsConfig. |
| 169 | + fn parse(yaml: &str) -> LogsConfig { |
| 170 | + serde_yaml::from_str(yaml).unwrap() |
| 171 | + } |
| 172 | + |
| 173 | + /// Helper to create LoggingProviders with specified modes. |
| 174 | + fn providers( |
| 175 | + global: ProviderMode, |
| 176 | + engine: ProviderMode, |
| 177 | + internal: ProviderMode, |
| 178 | + ) -> LoggingProviders { |
| 179 | + LoggingProviders { |
| 180 | + global, |
| 181 | + engine, |
| 182 | + internal, |
| 183 | + } |
| 184 | + } |
| 185 | + |
| 186 | + /// Helper to create a config with custom providers. |
| 187 | + fn config_with( |
| 188 | + global: ProviderMode, |
| 189 | + engine: ProviderMode, |
| 190 | + internal: ProviderMode, |
| 191 | + ) -> LogsConfig { |
| 192 | + LogsConfig { |
| 193 | + providers: providers(global, engine, internal), |
| 194 | + ..Default::default() |
| 195 | + } |
| 196 | + } |
| 197 | + |
| 198 | + /// Asserts validation fails with expected substring in error message. |
| 199 | + fn assert_invalid(config: &LogsConfig, expected_msg: &str) { |
| 200 | + let err = config.validate().unwrap_err(); |
| 201 | + assert!(matches!(err, Error::InvalidUserConfig { .. })); |
| 202 | + assert!( |
| 203 | + err.to_string().contains(expected_msg), |
| 204 | + "Expected '{}' in: {}", |
| 205 | + expected_msg, |
| 206 | + err |
| 207 | + ); |
| 208 | + } |
| 209 | + |
47 | 210 | #[test] |
48 | | - fn test_logs_config_deserialize() { |
49 | | - let yaml_str = r#" |
50 | | - level: "info" |
51 | | - processors: |
52 | | - - batch: |
53 | | - exporter: |
54 | | - console: |
55 | | - "#; |
56 | | - let config: LogsConfig = serde_yaml::from_str(yaml_str).unwrap(); |
| 211 | + fn test_defaults() { |
| 212 | + // Manual Default impl matches serde defaults |
| 213 | + let config = LogsConfig::default(); |
57 | 214 | assert_eq!(config.level, LogLevel::Info); |
58 | | - assert_eq!(config.processors.len(), 1); |
| 215 | + assert_eq!(config.providers.global, ProviderMode::ConsoleAsync); |
| 216 | + assert_eq!(config.providers.engine, ProviderMode::ConsoleAsync); |
| 217 | + assert_eq!(config.providers.internal, ProviderMode::Noop); |
| 218 | + assert!(config.processors.is_empty()); |
| 219 | + |
| 220 | + // Serde defaults should match Rust Default |
| 221 | + let parsed = parse("{}"); |
| 222 | + assert_eq!(parsed.level, config.level); |
| 223 | + assert_eq!(parsed.providers.global, config.providers.global); |
| 224 | + assert_eq!(parsed.providers.engine, config.providers.engine); |
| 225 | + assert_eq!(parsed.providers.internal, config.providers.internal); |
59 | 226 | } |
60 | 227 |
|
61 | 228 | #[test] |
62 | | - fn test_log_level_deserialize() { |
63 | | - let yaml_str = r#" |
64 | | - level: "info" |
65 | | - "#; |
66 | | - let config: LogsConfig = serde_yaml::from_str(yaml_str).unwrap(); |
67 | | - assert_eq!(config.level, LogLevel::Info); |
| 229 | + fn test_log_level_parsing() { |
| 230 | + let cases = [ |
| 231 | + ("off", LogLevel::Off), |
| 232 | + ("debug", LogLevel::Debug), |
| 233 | + ("info", LogLevel::Info), |
| 234 | + ("warn", LogLevel::Warn), |
| 235 | + ("error", LogLevel::Error), |
| 236 | + ]; |
| 237 | + for (name, expected) in cases { |
| 238 | + assert_eq!(parse(&format!("level: {name}")).level, expected); |
| 239 | + } |
68 | 240 | } |
69 | 241 |
|
70 | 242 | #[test] |
71 | | - fn test_logs_config_default_deserialize() -> Result<(), serde_yaml::Error> { |
72 | | - let yaml_str = r#""#; |
73 | | - let config: LogsConfig = serde_yaml::from_str(yaml_str)?; |
74 | | - assert_eq!(config.level, LogLevel::Off); |
75 | | - assert!(config.processors.is_empty()); |
76 | | - Ok(()) |
| 243 | + fn test_provider_mode_parsing() { |
| 244 | + let config = parse("providers: { global: noop, engine: its, internal: console_direct }"); |
| 245 | + assert_eq!(config.providers.global, ProviderMode::Noop); |
| 246 | + assert_eq!(config.providers.engine, ProviderMode::ITS); |
| 247 | + assert_eq!(config.providers.internal, ProviderMode::ConsoleDirect); |
| 248 | + |
| 249 | + let config = parse("providers: { global: opentelemetry, engine: opentelemetry }"); |
| 250 | + assert_eq!(config.providers.global, ProviderMode::OpenTelemetry); |
| 251 | + assert_eq!(config.providers.engine, ProviderMode::OpenTelemetry); |
| 252 | + } |
| 253 | + |
| 254 | + #[test] |
| 255 | + fn test_needs_reporter() { |
| 256 | + use ProviderMode::*; |
| 257 | + let cases = [ |
| 258 | + (Noop, false), |
| 259 | + (ITS, true), |
| 260 | + (OpenTelemetry, false), |
| 261 | + (ConsoleDirect, false), |
| 262 | + (ConsoleAsync, true), |
| 263 | + ]; |
| 264 | + for (mode, expected) in cases { |
| 265 | + assert_eq!(mode.needs_reporter(), expected, "{mode:?}"); |
| 266 | + } |
| 267 | + } |
| 268 | + |
| 269 | + #[test] |
| 270 | + fn test_validate_default_succeeds() { |
| 271 | + assert!(LogsConfig::default().validate().is_ok()); |
| 272 | + } |
| 273 | + |
| 274 | + #[test] |
| 275 | + fn test_validate_internal_cannot_use_reporter() { |
| 276 | + use ProviderMode::*; |
| 277 | + let config = config_with(Noop, Noop, ITS); |
| 278 | + assert_invalid(&config, "internal provider is invalid"); |
| 279 | + |
| 280 | + let config = config_with(Noop, Noop, ConsoleAsync); |
| 281 | + assert_invalid(&config, "internal provider is invalid"); |
| 282 | + } |
| 283 | + |
| 284 | + #[test] |
| 285 | + fn test_validate_engine_otel_requires_global_otel() { |
| 286 | + use ProviderMode::*; |
| 287 | + // Engine OpenTelemetry without global OpenTelemetry fails |
| 288 | + for global in [Noop, ITS, ConsoleDirect, ConsoleAsync] { |
| 289 | + let config = config_with(global, OpenTelemetry, Noop); |
| 290 | + assert_invalid(&config, "opentelemetry"); |
| 291 | + } |
| 292 | + |
| 293 | + // Both OpenTelemetry succeeds |
| 294 | + let config = config_with(OpenTelemetry, OpenTelemetry, Noop); |
| 295 | + assert!(config.validate().is_ok()); |
77 | 296 | } |
78 | 297 | } |
0 commit comments