diff --git a/Cargo.lock b/Cargo.lock index d32301556d..8672dec631 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7049,6 +7049,7 @@ version = "0.5.1-edge.1" dependencies = [ "axum", "axum-server", + "clap", "configs", "configs_derive", "dirs 7.0.0", diff --git a/core/ai/mcp/Cargo.toml b/core/ai/mcp/Cargo.toml index 1fe5afb44c..a912c0255f 100644 --- a/core/ai/mcp/Cargo.toml +++ b/core/ai/mcp/Cargo.toml @@ -33,6 +33,7 @@ systemd = ["dep:sd-notify", "dep:tokio-util"] [dependencies] axum = { workspace = true } axum-server = { workspace = true } +clap = { workspace = true } configs = { workspace = true } configs_derive = { workspace = true } dirs = { workspace = true } diff --git a/core/ai/mcp/README.md b/core/ai/mcp/README.md index 12ed06259a..3b9754a86a 100644 --- a/core/ai/mcp/README.md +++ b/core/ai/mcp/README.md @@ -58,10 +58,14 @@ The configuration file must use TOML. The default path is `core/ai/mcp/config.to Set `IGGY_MCP_ENV_PATH` to load a particular dotenv file. Otherwise `.env` is searched for in the current directory and its parents. Existing environment variables take precedence over dotenv values. +Run `iggy-mcp --list-config-env-vars` to print the supported configuration environment variables and exit before loading dotenv or configuration files or creating the runtime. + A non-empty `iggy.token` takes precedence over username and password. It accepts a literal PAT or a `file:` reference such as `file:/run/secrets/iggy_pat`; file contents are trimmed and a leading `~/` expands to the home directory. Set `command` to the absolute path of the built executable. This Claude Desktop example uses the development broker credentials: +`iggy-mcp` rejects unknown command-line arguments with exit code 2. Keep the client or container `args` list empty unless it contains a supported flag. + ```json { "mcpServers": { diff --git a/core/ai/mcp/src/error.rs b/core/ai/mcp/src/error.rs index 5b6e0347e1..369938bd3f 100644 --- a/core/ai/mcp/src/error.rs +++ b/core/ai/mcp/src/error.rs @@ -45,4 +45,6 @@ pub enum McpRuntimeError { TokenFileReadError(String, String), #[error("Token file is empty: {0}")] TokenFileEmpty(String), + #[error("Failed to list config environment variables")] + ListConfigEnvVars(#[source] std::io::Error), } diff --git a/core/ai/mcp/src/main.rs b/core/ai/mcp/src/main.rs index c2d24fbd26..38dded37da 100644 --- a/core/ai/mcp/src/main.rs +++ b/core/ai/mcp/src/main.rs @@ -15,7 +15,11 @@ // specific language governing permissions and limitations // under the License. -use ::configs::ConfigProvider; +use ::configs::{ + ConfigEnvMappings, ConfigProvider, MCP_CONFIG_PATH_ENV, MCP_ENV_PATH_ENV, MCP_RUNTIME_ENV_VARS, + print_env_var_names, +}; +use clap::Parser; use configs::{McpServerConfig, McpTransport}; use dotenvy::dotenv; use error::McpRuntimeError; @@ -42,7 +46,20 @@ const VERSION: &str = env!("CARGO_PKG_VERSION"); const DEFAULT_CONFIG_PATH: &str = "core/ai/mcp/config.toml"; +#[derive(Debug, Parser)] +#[command(author = "Apache Iggy", version)] +struct Args { + /// Print supported configuration environment variables and exit. + #[arg(long)] + list_config_env_vars: bool, +} + fn main() -> Result<(), McpRuntimeError> { + let args = Args::parse(); + if args.list_config_env_vars { + print_config_env_vars().map_err(McpRuntimeError::ListConfigEnvVars)?; + return Ok(()); + } let runtime = Builder::new_multi_thread() .enable_all() .build() @@ -53,12 +70,23 @@ fn main() -> Result<(), McpRuntimeError> { result } +fn print_config_env_vars() -> std::io::Result<()> { + let mut stdout = std::io::stdout(); + print_env_var_names( + McpServerConfig::env_templates() + .iter() + .map(|t| t.env_name) + .chain(MCP_RUNTIME_ENV_VARS.iter().copied()), + &mut stdout, + ) +} + async fn run() -> Result<(), McpRuntimeError> { let standard_font = FIGlet::standard().unwrap(); let figure = standard_font.convert("Iggy MCP Server"); eprintln!("{}", figure.unwrap()); - if let Ok(env_path) = std::env::var("IGGY_MCP_ENV_PATH") { + if let Ok(env_path) = std::env::var(MCP_ENV_PATH_ENV) { if dotenvy::from_path(&env_path).is_ok() { eprintln!("Loaded environment variables from path: {env_path}"); } @@ -70,7 +98,7 @@ async fn run() -> Result<(), McpRuntimeError> { } let config_path = - env::var("IGGY_MCP_CONFIG_PATH").unwrap_or_else(|_| DEFAULT_CONFIG_PATH.to_string()); + env::var(MCP_CONFIG_PATH_ENV).unwrap_or_else(|_| DEFAULT_CONFIG_PATH.to_string()); eprintln!("Configuration file path: {config_path}"); let config: McpServerConfig = McpServerConfig::config_provider(config_path) .load_config() diff --git a/core/configs/src/configs_impl/env_listing.rs b/core/configs/src/configs_impl/env_listing.rs new file mode 100644 index 0000000000..17512ab90d --- /dev/null +++ b/core/configs/src/configs_impl/env_listing.rs @@ -0,0 +1,102 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +/// Writes each name in `names`, sorted and deduplicated, to the provided writer — one +/// call, one error policy. A closed writer (e.g. `iggy-server --list-config-env-vars | head -1`) +/// is expected, not a failure: on `BrokenPipe` this stops writing and returns `Ok(())`. Any +/// other I/O error is real and is propagated so the caller exits non-zero. +pub fn print_env_var_names(names: I, writer: &mut W) -> std::io::Result<()> +where + I: IntoIterator, + S: Into, + W: std::io::Write, +{ + let mut names: Vec = names.into_iter().map(Into::into).collect(); + names.sort_unstable(); + names.dedup(); + + for name in names { + if let Err(err) = writeln!(writer, "{name}") { + return if err.kind() == std::io::ErrorKind::BrokenPipe { + Ok(()) + } else { + Err(err) + }; + } + } + Ok(()) +} + +pub const MCP_CONFIG_PATH_ENV: &str = "IGGY_MCP_CONFIG_PATH"; +pub const MCP_ENV_PATH_ENV: &str = "IGGY_MCP_ENV_PATH"; +/// Env vars `iggy-mcp --list-config-env-vars` advertises beyond the derived +/// `McpServerConfig` templates. +pub const MCP_RUNTIME_ENV_VARS: &[&str] = &[ + super::file_provider::DISPLAY_CONFIG_ENV, + MCP_CONFIG_PATH_ENV, + MCP_ENV_PATH_ENV, +]; + +pub const CONNECTORS_CONFIG_PATH_ENV: &str = "IGGY_CONNECTORS_CONFIG_PATH"; +pub const CONNECTORS_ENV_PATH_ENV: &str = "IGGY_CONNECTORS_ENV_PATH"; +/// Env vars `iggy-connectors --list-config-env-vars` advertises beyond the +/// derived `ConnectorsRuntimeConfig` templates. +pub const CONNECTORS_RUNTIME_ENV_VARS: &[&str] = &[ + CONNECTORS_CONFIG_PATH_ENV, + CONNECTORS_ENV_PATH_ENV, + super::file_provider::DISPLAY_CONFIG_ENV, +]; + +/// Segment between a connector env-var prefix and a plugin config field. +pub const PLUGIN_CONFIG_ENV_SEGMENT: &str = "PLUGIN_CONFIG_"; + +#[cfg(test)] +mod tests { + use super::*; + use std::io::{Error, ErrorKind, Write}; + + struct FailingWriter(ErrorKind); + + impl Write for FailingWriter { + fn write(&mut self, _: &[u8]) -> std::io::Result { + Err(Error::from(self.0)) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + #[test] + fn print_env_var_names_sorts_and_deduplicates() { + let mut output = Vec::new(); + print_env_var_names(["IGGY_B", "IGGY_A", "IGGY_B"], &mut output).unwrap(); + assert_eq!(output, b"IGGY_A\nIGGY_B\n"); + } + + #[test] + fn print_env_var_names_accepts_a_closed_pipe() { + assert!(print_env_var_names(["IGGY_A"], &mut FailingWriter(ErrorKind::BrokenPipe)).is_ok()); + } + + #[test] + fn print_env_var_names_propagates_other_write_errors() { + let error = print_env_var_names(["IGGY_A"], &mut FailingWriter(ErrorKind::StorageFull)) + .expect_err("storage-full error must be propagated"); + assert_eq!(error.kind(), ErrorKind::StorageFull); + } +} diff --git a/core/configs/src/configs_impl/env_mapping.rs b/core/configs/src/configs_impl/env_mapping.rs index 11336f7e27..673e9892fc 100644 --- a/core/configs/src/configs_impl/env_mapping.rs +++ b/core/configs/src/configs_impl/env_mapping.rs @@ -29,12 +29,31 @@ pub struct EnvVarMapping { pub is_secret: bool, } +/// Compact representation of a configuration environment variable. +/// +/// Array indices are represented by `` in `env_name`. `max_elements` +/// contains the expansion limit for each placeholder, from left to right. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct EnvVarTemplate { + /// Environment variable name, with `` for each array index. + pub env_name: &'static str, + /// Config path, with `` for each array index. + pub config_path: &'static str, + /// Whether this field contains secret data. + pub is_secret: bool, + /// Maximum element count for each `` placeholder, from left to right. + pub max_elements: &'static [usize], +} + /// Trait for configuration types that provide environment variable mappings. /// Implemented automatically by the `#[derive(ConfigEnv)]` macro. pub trait ConfigEnvMappings { /// Returns all environment variable mappings for this config type. fn env_mappings() -> &'static [EnvVarMapping]; + /// Returns the compact environment variable templates for this config type. + fn env_templates() -> &'static [EnvVarTemplate]; + /// Finds a mapping by environment variable name. fn find_by_env_name(env_name: &str) -> Option<&'static EnvVarMapping> { Self::env_mappings().iter().find(|m| m.env_name == env_name) @@ -59,3 +78,41 @@ pub trait ConfigEnvMappings { .collect() } } + +impl EnvVarTemplate { + fn expand(&self) -> Vec { + let mut results = vec![(self.env_name.to_string(), self.config_path.to_string())]; + + for &limit in self.max_elements { + let mut next = Vec::new(); + for (env_name, config_path) in results { + for i in 0..limit { + let index = i.to_string(); + next.push(( + env_name.replacen("", &index, 1), + config_path.replacen("", &index, 1), + )); + } + } + results = next; + } + + results + .into_iter() + .map(|(env_name, config_path)| EnvVarMapping { + env_name: Box::leak(env_name.into_boxed_str()), + config_path: Box::leak(config_path.into_boxed_str()), + is_secret: self.is_secret, + }) + .collect() + } +} + +/// Expands compact templates into the mappings consumed by the env provider. +/// +/// The returned mappings own leaked names and paths so generated +/// `ConfigEnvMappings` implementations can cache them in a `OnceLock` and +/// return `'static` references. Call this once per config type. +pub fn expand_env_templates(templates: &[EnvVarTemplate]) -> Vec { + templates.iter().flat_map(EnvVarTemplate::expand).collect() +} diff --git a/core/configs/src/configs_impl/file_provider.rs b/core/configs/src/configs_impl/file_provider.rs index 851000303a..33b844fb78 100644 --- a/core/configs/src/configs_impl/file_provider.rs +++ b/core/configs/src/configs_impl/file_provider.rs @@ -26,7 +26,7 @@ use figment::{ use std::{env, path::Path}; use tracing::{error, info, warn}; -const DISPLAY_CONFIG_ENV: &str = "IGGY_DISPLAY_CONFIG"; +pub(crate) const DISPLAY_CONFIG_ENV: &str = "IGGY_DISPLAY_CONFIG"; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum RelocatedTarget { @@ -448,7 +448,7 @@ mod tests { let unknown = unknown_env_names( names(&siblings).into_iter(), "IGGY_", - crate::server_config::server::SERVER_PROCESS_ENV_VARS, + &crate::server_config::server::server_process_env_vars().collect::>(), crate::server_config::server::SERVER_ALLOWED_ENV_PREFIXES, ); assert!( @@ -483,12 +483,13 @@ mod tests { } /// The server reads these variables outside its config, so the boot check - /// must accept them. Without the `SERVER_PROCESS_ENV_VARS` chain in + /// must accept them. Without the `server_process_env_vars()` chain in /// `ServerConfig::config_provider`, a debug build refuses to boot. #[test] fn given_the_server_process_variables_when_checking_then_the_server_should_boot() { - let unknown = - server_unknown_env_names(crate::server_config::server::SERVER_PROCESS_ENV_VARS); + let process_env_vars = + crate::server_config::server::server_process_env_vars().collect::>(); + let unknown = server_unknown_env_names(&process_env_vars); assert!( unknown.is_empty(), diff --git a/core/configs/src/configs_impl/mod.rs b/core/configs/src/configs_impl/mod.rs index b94d4ddd63..a9bf2bf91c 100644 --- a/core/configs/src/configs_impl/mod.rs +++ b/core/configs/src/configs_impl/mod.rs @@ -24,6 +24,7 @@ //! - JSON value field handling //! - Automatic type conversion and validation +mod env_listing; mod env_mapping; mod error; mod file_provider; @@ -31,7 +32,12 @@ mod parsing; mod traits; mod typed_env_provider; -pub use env_mapping::{ConfigEnvMappings, EnvVarMapping}; +pub use env_listing::{ + CONNECTORS_CONFIG_PATH_ENV, CONNECTORS_ENV_PATH_ENV, CONNECTORS_RUNTIME_ENV_VARS, + MCP_CONFIG_PATH_ENV, MCP_ENV_PATH_ENV, MCP_RUNTIME_ENV_VARS, PLUGIN_CONFIG_ENV_SEGMENT, + print_env_var_names, +}; +pub use env_mapping::{ConfigEnvMappings, EnvVarMapping, EnvVarTemplate, expand_env_templates}; pub use error::ConfigurationError; pub use file_provider::{FileConfigProvider, RelocatedKey, RelocatedTarget}; pub use parsing::parse_env_value_to_json; diff --git a/core/configs/src/configs_impl/typed_env_provider.rs b/core/configs/src/configs_impl/typed_env_provider.rs index b6be617da5..74dd06324b 100644 --- a/core/configs/src/configs_impl/typed_env_provider.rs +++ b/core/configs/src/configs_impl/typed_env_provider.rs @@ -15,6 +15,7 @@ // specific language governing permissions and limitations // under the License. +use super::PLUGIN_CONFIG_ENV_SEGMENT; use super::env_mapping::ConfigEnvMappings; use super::error::ConfigurationError; use super::parsing::parse_env_value; @@ -37,26 +38,21 @@ enum EnvNameResolution<'a> { /// Controls filtering and messaging for unknown env var warnings. enum WarningContext<'a> { - /// Main config: filter IGNORED_ENV_VARS and DELEGATED prefixes + /// Main config: filter runtime env vars and delegated prefixes. MainConfig, /// Connector config: filter PLUGIN_CONFIG_ prefix ConnectorConfig(&'a str), } -/// `IGGY_` variables that are NOT config values: the config file and dotenv -/// paths the connectors runtime and the MCP server read before their config -/// loads. -const IGNORED_ENV_VARS: &[&str] = &[ - "IGGY_CONNECTORS_CONFIG_PATH", - "IGGY_CONNECTORS_ENV_PATH", - "IGGY_MCP_CONFIG_PATH", - "IGGY_MCP_ENV_PATH", -]; - /// Prefixes for env vars handled by separate providers with runtime prefixes. /// The main config provider skips these; each sub-provider validates its own vars. const DELEGATED_ENV_VAR_PREFIXES: &[&str] = &["IGGY_CONNECTORS_SINK_", "IGGY_CONNECTORS_SOURCE_"]; +fn is_runtime_env_var(name: &str) -> bool { + super::CONNECTORS_RUNTIME_ENV_VARS.contains(&name) + || super::MCP_RUNTIME_ENV_VARS.contains(&name) +} + type ProfileMap = FigmentMap; /// Type-safe environment variable provider that uses compile-time generated mappings. @@ -291,13 +287,13 @@ impl TypedEnvProvider { let should_skip = match &context { WarningContext::MainConfig => { - IGNORED_ENV_VARS.contains(&key.as_str()) + is_runtime_env_var(&key) || DELEGATED_ENV_VAR_PREFIXES .iter() .any(|p| key.starts_with(p)) } WarningContext::ConnectorConfig(prefix) => { - let plugin_config_prefix = format!("{}PLUGIN_CONFIG_", prefix); + let plugin_config_prefix = format!("{}{PLUGIN_CONFIG_ENV_SEGMENT}", prefix); key.starts_with(&plugin_config_prefix) } }; @@ -314,7 +310,7 @@ impl TypedEnvProvider { let debug_msg = match &context { WarningContext::MainConfig => format!( - "Unknown IGGY_ env var: '{}'.{}. Add to IGNORED_ENV_VARS if intentional, \ + "Unknown IGGY_ env var: '{}'.{}. Add to the runtime env-var lists if intentional, \ or add #[derive(ConfigEnv)] to the config struct.", key, suggestion_hint ), @@ -625,19 +621,21 @@ mod tests { #[test] #[serial_test::serial] fn ignored_env_vars_are_skipped_by_the_unknown_variable_scan() { - for name in [ + let runtime_env_vars = [ "IGGY_CONNECTORS_CONFIG_PATH", "IGGY_CONNECTORS_ENV_PATH", "IGGY_MCP_CONFIG_PATH", "IGGY_MCP_ENV_PATH", - ] { + ]; + + for name in runtime_env_vars { assert!( - IGNORED_ENV_VARS.contains(&name), - "{name} is read by a sibling binary before its config loads, so the scan must skip it" + is_runtime_env_var(name), + "{name} is read before config loading and must be skipped by the scan" ); } - for name in IGNORED_ENV_VARS { + for name in runtime_env_vars { // SAFETY: the race is process-wide, not per key: `set_var` is unsound // against any concurrent environment access. `serial_test::serial` on // this test is what prevents that. @@ -649,7 +647,7 @@ mod tests { .warn_unknown_env_vars_inner(WarningContext::MainConfig); } - for name in IGNORED_ENV_VARS { + for name in runtime_env_vars { // SAFETY: paired with the set above. unsafe { env::remove_var(name) }; } diff --git a/core/configs/src/lib.rs b/core/configs/src/lib.rs index 62832e1362..8e266121bb 100644 --- a/core/configs/src/lib.rs +++ b/core/configs/src/lib.rs @@ -23,8 +23,11 @@ mod server_config; pub use common::{COMPONENT, defaults, displays, http, system, validators}; pub use configs_derive::ConfigEnv; pub use configs_impl::{ + CONNECTORS_CONFIG_PATH_ENV, CONNECTORS_ENV_PATH_ENV, CONNECTORS_RUNTIME_ENV_VARS, ConfigEnvMappings, ConfigProvider, ConfigurationError, ConfigurationType, EnvVarMapping, - FileConfigProvider, RelocatedKey, RelocatedTarget, TypedEnvProvider, parse_env_value_to_json, + EnvVarTemplate, FileConfigProvider, MCP_CONFIG_PATH_ENV, MCP_ENV_PATH_ENV, + MCP_RUNTIME_ENV_VARS, PLUGIN_CONFIG_ENV_SEGMENT, RelocatedKey, RelocatedTarget, + TypedEnvProvider, expand_env_templates, parse_env_value_to_json, print_env_var_names, }; pub use server_config::{ cluster, message_bus, metadata, partition, quic, server, sharding, tcp, websocket, diff --git a/core/configs/src/server_config/cluster.rs b/core/configs/src/server_config/cluster.rs index 0667c131cc..140847951a 100644 --- a/core/configs/src/server_config/cluster.rs +++ b/core/configs/src/server_config/cluster.rs @@ -1569,6 +1569,25 @@ mod tests { "selector env expansion must stop at max_elements = 16" ); } + + #[test] + fn advertised_addresses_env_template_keeps_each_vector_limit() { + let template = ::env_templates() + .iter() + .find(|template| template.env_name == "NODES__ADVERTISED_ADDRESSES__CLIENT_CIDR") + .expect("advertised address template"); + + assert_eq!(template.max_elements, &[256, 16]); + + let mapping = ::env_mappings() + .iter() + .find(|mapping| mapping.env_name == "NODES_3_ADVERTISED_ADDRESSES_7_CLIENT_CIDR") + .expect("expanded advertised address mapping"); + assert_eq!( + mapping.config_path, + "nodes.3.advertised_addresses.7.client_cidr" + ); + } } #[cfg(test)] diff --git a/core/configs/src/server_config/server.rs b/core/configs/src/server_config/server.rs index 9f852815fa..5963203aa0 100644 --- a/core/configs/src/server_config/server.rs +++ b/core/configs/src/server_config/server.rs @@ -47,23 +47,42 @@ pub use crate::common::server::{ TelemetryConfig, TelemetryLogsConfig, TelemetryTracesConfig, TelemetryTransport, }; -pub const SERVER_PROCESS_ENV_VARS: &[&str] = &[ - "IGGY_CONFIG_PATH", - "IGGY_ENV_PATH", - "IGGY_DISPLAY_CONFIG", - "IGGY_ROOT_USERNAME", - "IGGY_ROOT_PASSWORD", +/// Vars used by sibling binaries (iggy CLI) or test/CI-only. These suppress +/// unknown-name warnings but are not advertised by the server. +const SERVER_SCAN_ONLY_ENV_VARS: &[&str] = &[ "IGGY_TEST_VERBOSE", "IGGY_TEST_CLUSTER_NODES", "IGGY_TEST_CLEANUP_DISABLED", - "IGGY_SHARD_RUNTIME_CAPACITY", - "IGGY_SHARD_EVENT_INTERVAL", "IGGY_CI_BUILD", "IGGY_HOME", "IGGY_USERNAME", "IGGY_PASSWORD", ]; +/// Non-config env vars supported by the server and advertised to operators. +const SERVER_RUNTIME_ENV_VARS: &[&str] = &[ + "IGGY_CONFIG_PATH", + "IGGY_ENV_PATH", + "IGGY_DISPLAY_CONFIG", + "IGGY_ROOT_USERNAME", + "IGGY_ROOT_PASSWORD", + "IGGY_SHARD_RUNTIME_CAPACITY", + "IGGY_SHARD_EVENT_INTERVAL", +]; + +/// Non-config vars advertised by `iggy-server --list-config-env-vars`. +pub fn server_runtime_env_vars() -> impl Iterator { + SERVER_RUNTIME_ENV_VARS.iter().copied() +} + +/// All non-config vars accepted by the server's unknown-name scan. +pub fn server_process_env_vars() -> impl Iterator { + SERVER_RUNTIME_ENV_VARS + .iter() + .chain(SERVER_SCAN_ONLY_ENV_VARS) + .copied() +} + pub(crate) const SERVER_ALLOWED_ENV_PREFIXES: &[&str] = &["IGGY_CONNECTORS_", "IGGY_KAFKA_", "IGGY_MCP_"]; @@ -262,7 +281,7 @@ impl ServerConfig { .with_known_env_names( Self::all_env_var_names() .into_iter() - .chain(SERVER_PROCESS_ENV_VARS.iter().copied()) + .chain(server_process_env_vars()) .collect(), ) .with_allowed_env_prefixes(SERVER_ALLOWED_ENV_PREFIXES) @@ -468,7 +487,7 @@ mod tests { #[test] #[serial_test::serial] fn env_provider_accepts_server_process_env_vars() { - for name in SERVER_PROCESS_ENV_VARS { + for name in server_process_env_vars() { // SAFETY: the race is process-wide, not per key: `set_var` is unsound // against any concurrent environment access. `serial_test::serial` on // this test is what prevents that. @@ -477,7 +496,7 @@ mod tests { let data = ServerConfigEnvProvider::default().data(); - for name in SERVER_PROCESS_ENV_VARS { + for name in server_process_env_vars() { // SAFETY: paired with the set above. unsafe { env::remove_var(name) }; } @@ -492,4 +511,29 @@ mod tests { "none of these variables is a config value, so none of them may reach the map: {profile:?}" ); } + + #[test] + fn server_process_env_vars_include_every_scan_only_name() { + let process_names = server_process_env_vars().collect::>(); + for name in [ + "IGGY_TEST_VERBOSE", + "IGGY_TEST_CLUSTER_NODES", + "IGGY_TEST_CLEANUP_DISABLED", + "IGGY_CI_BUILD", + "IGGY_HOME", + "IGGY_USERNAME", + "IGGY_PASSWORD", + ] { + assert!( + process_names.contains(name), + "missing scan-only name {name}" + ); + } + assert!( + SERVER_RUNTIME_ENV_VARS + .iter() + .all(|name| !SERVER_SCAN_ONLY_ENV_VARS.contains(name)), + "advertised and scan-only env vars must stay disjoint" + ); + } } diff --git a/core/configs_derive/src/config_env.rs b/core/configs_derive/src/config_env.rs index 69dbc0c85f..176efe0f02 100644 --- a/core/configs_derive/src/config_env.rs +++ b/core/configs_derive/src/config_env.rs @@ -180,46 +180,52 @@ fn generate_enum_impl( }) .collect(); - let tag_mapping = tag.map(|tag_name| { - let env_segment = tag_name.to_uppercase(); - quote! { - all_mappings.push(configs::EnvVarMapping { - env_name: #env_segment, - config_path: #tag_name, - is_secret: false, - }); - } - }); - - if variant_types.is_empty() && tag_mapping.is_none() { + if variant_types.is_empty() && tag.is_none() { return quote! { impl #impl_generics configs::ConfigEnvMappings for #enum_name #ty_generics #where_clause { fn env_mappings() -> &'static [configs::EnvVarMapping] { &[] } + + fn env_templates() -> &'static [configs::EnvVarTemplate] { + &[] + } } }; } - // Generate code to extend from each variant type - let extends: Vec = variant_types + let template_extends: Vec = variant_types .iter() - .map(|ty| { - quote! { - all_mappings.extend_from_slice(<#ty as configs::ConfigEnvMappings>::env_mappings()); - } + .map(|ty| quote! { + all_templates.extend_from_slice(<#ty as configs::ConfigEnvMappings>::env_templates()); }) .collect(); + let tag_template = tag.map(|tag_name| { + let env_segment = tag_name.to_uppercase(); + quote! { + all_templates.push(configs::EnvVarTemplate { + env_name: #env_segment, + config_path: #tag_name, + is_secret: false, + max_elements: &[], + }); + } + }); quote! { impl #impl_generics configs::ConfigEnvMappings for #enum_name #ty_generics #where_clause { fn env_mappings() -> &'static [configs::EnvVarMapping] { static MAPPINGS: std::sync::OnceLock> = std::sync::OnceLock::new(); - MAPPINGS.get_or_init(|| { - let mut all_mappings: Vec = Vec::new(); - #tag_mapping - #(#extends)* - all_mappings + MAPPINGS.get_or_init(|| configs::expand_env_templates(Self::env_templates())) + } + + fn env_templates() -> &'static [configs::EnvVarTemplate] { + static TEMPLATES: std::sync::OnceLock> = std::sync::OnceLock::new(); + TEMPLATES.get_or_init(|| { + let mut all_templates = Vec::new(); + #tag_template + #(#template_extends)* + all_templates }) } } @@ -237,8 +243,6 @@ fn generate_struct_impl( fields: Vec, ) -> TokenStream2 { let prefix_str = prefix.as_deref().unwrap_or(""); - let has_prefix = !prefix_str.is_empty(); - // Metadata name: use provided or derive from type name (e.g., "ServerConfig" -> "server-config") let metadata_name = name.unwrap_or_else(|| { let type_name = struct_name.to_string(); @@ -255,109 +259,59 @@ fn generate_struct_impl( let (mappings, nested_fields) = collect_mappings(&fields); let const_definitions = generate_const_definitions(&mappings, prefix_str); - let mapping_entries = generate_mapping_entries(&mappings); let builder_methods = generate_builder_methods(&mappings, prefix_str); let builder_name = format_ident!("{}EnvBuilder", struct_name); - let mappings_count = mappings.len(); - - // Generate code to include nested type mappings - let nested_extends: Vec = nested_fields - .iter() - .map(|info| { - let ty = &info.element_type; - let field_env_segment = &info.field_env_segment; - let field_name = &info.field_name; - - if info.is_vec { - // For Vec fields, expand with array indices - let max_elements = info.max_elements; - quote! { - let nested_mappings = <#ty as configs::ConfigEnvMappings>::env_mappings(); - for i in 0..#max_elements { - for mapping in nested_mappings { - // Transform env_name by prepending FIELD_INDEX_ - let new_env_name: &'static str = Box::leak( - format!("{}_{}_{}", #field_env_segment, i, mapping.env_name).into_boxed_str() - ); - // Transform config_path by prepending field.index. - let new_config_path: &'static str = Box::leak( - format!("{}.{}.{}", #field_name, i, mapping.config_path).into_boxed_str() - ); - all_mappings.push(configs::EnvVarMapping { - env_name: new_env_name, - config_path: new_config_path, - is_secret: mapping.is_secret, - }); - } - } - } - } else { - // For non-Vec nested types, transform by prepending field name - quote! { - let nested_mappings = <#ty as configs::ConfigEnvMappings>::env_mappings(); - for mapping in nested_mappings { - // Transform env_name by prepending FIELD_ - let new_env_name: &'static str = Box::leak( - format!("{}_{}", #field_env_segment, mapping.env_name).into_boxed_str() - ); - // Transform config_path by prepending field. - let new_config_path: &'static str = Box::leak( - format!("{}.{}", #field_name, mapping.config_path).into_boxed_str() - ); - all_mappings.push(configs::EnvVarMapping { - env_name: new_env_name, - config_path: new_config_path, - is_secret: mapping.is_secret, - }); - } - } - } - }) - .collect(); - - let prefix_application = if has_prefix { - quote! { - let all_mappings: Vec = all_mappings - .into_iter() - .map(|m| configs::EnvVarMapping { - env_name: Box::leak(format!("{}{}", #prefix_str, m.env_name).into_boxed_str()), - config_path: m.config_path, - is_secret: m.is_secret, - }) - .collect(); - } - } else { - quote! {} - }; - // If there are nested types, we need to use OnceLock to combine mappings - let env_mappings_impl = if nested_fields.is_empty() && !has_prefix { - // Simple case: no nested fields, no prefix transformation needed + let own_template_entries = mappings.iter().map(|mapping| { + let env_name = format!("{}{}", prefix_str, mapping.env_suffix); + let config_path = &mapping.config_path; + let is_secret = mapping.is_secret; quote! { - fn env_mappings() -> &'static [configs::EnvVarMapping] { - static MAPPINGS: [configs::EnvVarMapping; #mappings_count] = [ - #(#mapping_entries),* - ]; - &MAPPINGS + configs::EnvVarTemplate { + env_name: #env_name, + config_path: #config_path, + is_secret: #is_secret, + max_elements: &[], } } - } else { - // Complex case: need OnceLock for dynamic construction + }); + let nested_template_extends = nested_fields.iter().map(|info| { + let ty = &info.element_type; + let env_segment = if info.is_vec { + format!("{}{}_", prefix_str, info.field_env_segment) + } else { + format!("{}{}", prefix_str, info.field_env_segment) + }; + let config_segment = if info.is_vec { + format!("{}.", info.field_name) + } else { + info.field_name.clone() + }; + let limits = if info.is_vec { + let max_elements = info.max_elements; + quote! { + Box::leak(std::iter::once(#max_elements) + .chain(template.max_elements.iter().copied()) + .collect::>() + .into_boxed_slice()) + } + } else { + quote! { template.max_elements } + }; quote! { - fn env_mappings() -> &'static [configs::EnvVarMapping] { - static MAPPINGS: std::sync::OnceLock> = std::sync::OnceLock::new(); - MAPPINGS.get_or_init(|| { - let own_mappings: [configs::EnvVarMapping; #mappings_count] = [ - #(#mapping_entries),* - ]; - let mut all_mappings = Vec::from(own_mappings); - #(#nested_extends)* - #prefix_application - all_mappings - }) + for template in <#ty as configs::ConfigEnvMappings>::env_templates() { + let env_name = Box::leak(format!("{}_{}", #env_segment, template.env_name).into_boxed_str()); + let config_path = Box::leak(format!("{}.{}", #config_segment, template.config_path).into_boxed_str()); + let max_elements = #limits; + all_templates.push(configs::EnvVarTemplate { + env_name, + config_path, + is_secret: template.is_secret, + max_elements, + }); } } - }; + }); quote! { impl #impl_generics #struct_name #ty_generics #where_clause { @@ -371,7 +325,19 @@ fn generate_struct_impl( } impl #impl_generics configs::ConfigEnvMappings for #struct_name #ty_generics #where_clause { - #env_mappings_impl + fn env_mappings() -> &'static [configs::EnvVarMapping] { + static MAPPINGS: std::sync::OnceLock> = std::sync::OnceLock::new(); + MAPPINGS.get_or_init(|| configs::expand_env_templates(Self::env_templates())) + } + + fn env_templates() -> &'static [configs::EnvVarTemplate] { + static TEMPLATES: std::sync::OnceLock> = std::sync::OnceLock::new(); + TEMPLATES.get_or_init(|| { + let mut all_templates = vec![#(#own_template_entries),*]; + #(#nested_template_extends)* + all_templates + }) + } } /// Type-safe builder for constructing environment variable maps for tests. @@ -563,25 +529,6 @@ fn generate_const_definitions(mappings: &[EnvMapping], prefix: &str) -> Vec Vec { - mappings - .iter() - .map(|m| { - let env_suffix = &m.env_suffix; - let config_path = &m.config_path; - let is_secret = m.is_secret; - - quote! { - configs::EnvVarMapping { - env_name: #env_suffix, - config_path: #config_path, - is_secret: #is_secret, - } - } - }) - .collect() -} - fn generate_builder_methods(mappings: &[EnvMapping], prefix: &str) -> Vec { mappings .iter() diff --git a/core/connectors/runtime/README.md b/core/connectors/runtime/README.md index 580b7793eb..bf4f920e07 100644 --- a/core/connectors/runtime/README.md +++ b/core/connectors/runtime/README.md @@ -47,6 +47,10 @@ IGGY_CONNECTORS_CONFIG_PATH=connectors.toml cargo run --bin iggy-connectors Supported scalar fields and indexed list entries use environment variables with nested keys joined by underscores, for example `IGGY_CONNECTORS_IGGY_USERNAME`. Header and URL-template maps are configured in TOML. The runtime loads the first `.env` file found in the working directory or its parents, or the file specified by `IGGY_CONNECTORS_ENV_PATH`. +Run `iggy-connectors --list-config-env-vars` to print the supported names and templates without loading configuration or plugins. `` is a stream-vector index from 0 through 255. `` is the uppercased connector key and applies only to the local provider. + +`` sets one lowercased top-level `plugin_config` key with underscores preserved (`A_B` becomes `a_b`, not nested `a.b`); nested plugin configuration keys cannot be set through environment variables. `FORMAT` is reserved and listed separately as `..._PLUGIN_CONFIG_FORMAT`. + Source destination topics must persist every acknowledged batch before the runtime checkpoints the source or invokes its Ack hook. Missing topics are therefore created with `durability = "persisted"` and `messages_required_to_save = 1`. An existing topic must use `durability = "persisted"`; its save threshold may differ because persisted acknowledgments already wait for durable storage. ## State storage diff --git a/core/connectors/runtime/src/configs/connectors.rs b/core/connectors/runtime/src/configs/connectors.rs index ad6a6c435e..278481f4cc 100644 --- a/core/connectors/runtime/src/configs/connectors.rs +++ b/core/connectors/runtime/src/configs/connectors.rs @@ -16,7 +16,7 @@ // under the License. pub mod http_provider; -mod local_provider; +pub(crate) mod local_provider; use crate::configs::connectors::http_provider::HttpConnectorsConfigProvider; use crate::configs::connectors::local_provider::LocalConnectorsConfigProvider; diff --git a/core/connectors/runtime/src/configs/connectors/local_provider.rs b/core/connectors/runtime/src/configs/connectors/local_provider.rs index fc012bad25..e3c59791c6 100644 --- a/core/connectors/runtime/src/configs/connectors/local_provider.rs +++ b/core/connectors/runtime/src/configs/connectors/local_provider.rs @@ -21,7 +21,7 @@ use crate::configs::connectors::{ SourceConfig, }; use crate::error::RuntimeError; -use ::configs::{ConfigProvider, FileConfigProvider, TypedEnvProvider}; +use ::configs::{ConfigProvider, FileConfigProvider, PLUGIN_CONFIG_ENV_SEGMENT, TypedEnvProvider}; use async_trait::async_trait; use dashmap::DashMap; use figment::value::Dict; @@ -348,7 +348,10 @@ impl LocalConnectorsConfigProvider { ) { let connector_type = base_config.connector_type().to_uppercase(); let key = base_config.key().to_uppercase(); - let prefix = format!("IGGY_CONNECTORS_{connector_type}_{key}_PLUGIN_CONFIG_"); + let prefix = format!( + "{}{PLUGIN_CONFIG_ENV_SEGMENT}", + connector_env_prefix(&connector_type, &key) + ); for (env_key, env_value) in std::env::vars() { let env_key_upper = env_key.to_uppercase(); @@ -395,6 +398,20 @@ impl BaseConnectorConfig { } } +/// Builds the env-var prefix a connector's config and plugin-config overrides +/// are matched against: `IGGY_CONNECTORS_{TYPE}_{KEY}_`. `connector_type` and +/// `key` are expected uppercased, as `BaseConnectorConfig`'s callers already +/// do. Shared by the runtime's actual override lookup here and by +/// `--list-config-env-vars`'s listing in `main.rs` (passing the literal +/// `` as the key), so the two can't drift apart on the prefix the +/// connectors runtime actually reads. +pub(crate) fn connector_env_prefix(connector_type: &str, key: &str) -> String { + format!( + "{}{connector_type}_{key}_", + crate::configs::runtime::ConnectorsRuntimeConfig::ENV_PREFIX + ) +} + #[async_trait] impl ConnectorsConfigProvider for LocalConnectorsConfigProvider { async fn create_sink_config( @@ -822,7 +839,7 @@ impl ConnectorEnvProvider { fn with_connector_base_config(base_config: &BaseConnectorConfig) -> Self { let connector_type = base_config.connector_type().to_uppercase(); let key = base_config.key().to_uppercase(); - let prefix = format!("IGGY_CONNECTORS_{}_{}_", connector_type, key); + let prefix = connector_env_prefix(&connector_type, &key); let connector_name = base_config.key().to_owned(); match base_config { diff --git a/core/connectors/runtime/src/error.rs b/core/connectors/runtime/src/error.rs index 156427b131..9d415a7e6f 100644 --- a/core/connectors/runtime/src/error.rs +++ b/core/connectors/runtime/src/error.rs @@ -69,6 +69,10 @@ pub enum RuntimeError { CannotConvertConfiguration, #[error("IO operation failed with error: {0:?}")] IoError(#[from] std::io::Error), + #[error("Failed to list configuration environment variables")] + ListConfigEnvVars(#[source] std::io::Error), + #[error("Failed to create Tokio runtime")] + RuntimeCreation(#[source] std::io::Error), #[error("HTTP request failed: {0}")] HttpRequestFailed(String), #[error("Token file not found: {0}")] diff --git a/core/connectors/runtime/src/main.rs b/core/connectors/runtime/src/main.rs index 8890c1d430..ce0adf89c4 100644 --- a/core/connectors/runtime/src/main.rs +++ b/core/connectors/runtime/src/main.rs @@ -16,10 +16,14 @@ // under the License. use crate::configs::connectors::{ - ConnectorKey, ConnectorsConfig, ConnectorsConfigProvider, create_connectors_config_provider, + ConnectorKey, ConnectorsConfig, ConnectorsConfigProvider, SinkConfig, SourceConfig, + create_connectors_config_provider, }; use crate::metrics::ConnectorType; -use ::configs::ConfigProvider; +use ::configs::{ + CONNECTORS_CONFIG_PATH_ENV, CONNECTORS_ENV_PATH_ENV, CONNECTORS_RUNTIME_ENV_VARS, + ConfigEnvMappings, ConfigProvider, PLUGIN_CONFIG_ENV_SEGMENT, print_env_var_names, +}; use clap::Parser; use configs::connectors::ConfigFormat; use configs::runtime::ConnectorsRuntimeConfig; @@ -43,6 +47,7 @@ use std::{ sync::{Arc, atomic::AtomicU32}, }; use system_stats::capture_allowed_cpus; +use tokio::runtime::Builder; use tracing::{error, info, warn}; mod api; @@ -67,7 +72,26 @@ static GLOBAL: MiMalloc = MiMalloc; #[derive(Parser, Debug)] #[command(author = "Apache Iggy", version)] -struct Args {} +struct Args { + #[arg( + long, + help = "Print supported configuration environment variables and exit", + long_help = r#"Print supported configuration environment variables and exit. + +Lists all supported IGGY_* environment variable names and templates, +sorted and deduplicated. Template syntax: +- represents vector indices (0-255 for stream fields) +- represents connector keys (uppercased from config). + Overrides via require the local connectors provider. +- represents plugin configuration field names, excluding + FORMAT (handled separately as a strongly-typed field). It sets one + lowercased top-level plugin_config key: A_B becomes a_b, not nested a.b. + Nested plugin configuration keys cannot be set via environment variables. + +Exits immediately before any startup."# + )] + list_config_env_vars: bool, +} static PLUGIN_ID: AtomicU32 = AtomicU32::new(1); const ALLOWED_PLUGIN_EXTENSIONS: [&str; 3] = ["so", "dylib", "dll"]; @@ -117,13 +141,53 @@ fn print_ascii_art(text: &str) { println!("{}", figure.unwrap()); } -#[tokio::main] -async fn main() -> Result<(), RuntimeError> { +fn main() -> Result<(), RuntimeError> { + let args = Args::parse(); + if args.list_config_env_vars { + print_config_env_vars().map_err(RuntimeError::ListConfigEnvVars)?; + return Ok(()); + } capture_allowed_cpus(); - Args::parse(); + Builder::new_multi_thread() + .enable_all() + .build() + .map_err(RuntimeError::RuntimeCreation)? + .block_on(run()) +} + +fn print_config_env_vars() -> std::io::Result<()> { + let sink_source_templates = [ + ("SINK", SinkConfig::env_templates()), + ("SOURCE", SourceConfig::env_templates()), + ] + .into_iter() + .flat_map(|(kind, templates)| { + // "" stands in for the real, per-connector key `local_provider` + // uppercases at runtime - same prefix rule, so the listing can't + // drift from the names the runtime actually reads. + let prefix = + crate::configs::connectors::local_provider::connector_env_prefix(kind, ""); + let plugin_config = format!("{prefix}{PLUGIN_CONFIG_ENV_SEGMENT}"); + templates + .iter() + .map(move |template| format!("{prefix}{}", template.env_name)) + .chain(std::iter::once(plugin_config)) + }); + + let names = ConnectorsRuntimeConfig::env_templates() + .iter() + .map(|t| t.env_name.to_string()) + .chain(CONNECTORS_RUNTIME_ENV_VARS.iter().map(|s| s.to_string())) + .chain(sink_source_templates); + + let mut stdout = std::io::stdout(); + print_env_var_names(names, &mut stdout) +} + +async fn run() -> Result<(), RuntimeError> { print_ascii_art("Iggy Connectors"); - if let Ok(env_path) = std::env::var("IGGY_CONNECTORS_ENV_PATH") { + if let Ok(env_path) = std::env::var(CONNECTORS_ENV_PATH_ENV) { if dotenvy::from_path(&env_path).is_ok() { println!("Loaded environment variables from path: {env_path}"); } @@ -135,7 +199,7 @@ async fn main() -> Result<(), RuntimeError> { } let config_path = - env::var("IGGY_CONNECTORS_CONFIG_PATH").unwrap_or_else(|_| DEFAULT_CONFIG_PATH.to_string()); + env::var(CONNECTORS_CONFIG_PATH_ENV).unwrap_or_else(|_| DEFAULT_CONFIG_PATH.to_string()); println!("Starting Iggy Connectors Runtime, loading configuration from: {config_path}..."); let config: ConnectorsRuntimeConfig = ConnectorsRuntimeConfig::config_provider(config_path) diff --git a/core/integration/src/harness/config/resolve.rs b/core/integration/src/harness/config/resolve.rs index 10994ce372..1a50b83747 100644 --- a/core/integration/src/harness/config/resolve.rs +++ b/core/integration/src/harness/config/resolve.rs @@ -25,7 +25,9 @@ use std::collections::HashMap; /// `ServerConfig::all_env_var_names` cannot know them. `IGGY_CONFIG_PATH` /// selects the config file itself and the root credentials are consumed by /// `args.rs` before the config loads; `IGGY_TEST_VERBOSE` is harness-only. -pub const NON_CONFIG_ENV_VARS: &[&str] = configs::server::SERVER_PROCESS_ENV_VARS; +fn non_config_env_vars() -> impl Iterator { + configs::server::server_process_env_vars() +} /// Resolve config paths to environment variable names. /// @@ -111,8 +113,8 @@ fn find_mapping(path: &str) -> Option<&'static EnvVarMapping> { /// /// Names outside the `IGGY_` prefix are left alone: those address the process /// environment (`RUST_LOG`, test scaffolding), not the config schema. -/// `NON_CONFIG_ENV_VARS` carries the `IGGY_`-prefixed names the server reads -/// outside the config struct. +/// `server_process_env_vars()` supplies the `IGGY_`-prefixed names the server +/// reads outside the config struct. /// /// # Errors /// @@ -125,7 +127,7 @@ pub fn validate_env_var_names(envs: &HashMap) -> Result<(), Stri .filter(|name| { name.starts_with("IGGY_") && !known.contains(&name.as_str()) - && !NON_CONFIG_ENV_VARS.contains(&name.as_str()) + && !non_config_env_vars().any(|known| known == name.as_str()) }) .collect(); if unknown.is_empty() { @@ -283,9 +285,8 @@ mod tests { #[test] fn validate_env_var_names_accepts_the_non_config_variables() { - let envs: HashMap = NON_CONFIG_ENV_VARS - .iter() - .map(|name| ((*name).to_string(), "value".to_string())) + let envs: HashMap = non_config_env_vars() + .map(|name| (name.to_string(), "value".to_string())) .collect(); assert!( validate_env_var_names(&envs).is_ok(), diff --git a/core/integration/tests/config_env_listing/mod.rs b/core/integration/tests/config_env_listing/mod.rs new file mode 100644 index 0000000000..205b295e93 --- /dev/null +++ b/core/integration/tests/config_env_listing/mod.rs @@ -0,0 +1,294 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use assert_cmd::Command; +use std::collections::HashSet; +use std::time::Duration; + +const LIST_ENV_VARS_TIMEOUT: Duration = Duration::from_secs(5); + +fn list_config_env_vars_in( + binary: &str, + directory: &std::path::Path, + env: I, + extra_args: &[&str], +) -> Vec +where + I: IntoIterator, + K: AsRef, + V: AsRef, +{ + let mut cmd = Command::cargo_bin(binary).expect("binary should be built"); + cmd.current_dir(directory) + .arg("--list-config-env-vars") + .envs(env) + .args(extra_args) + .timeout(LIST_ENV_VARS_TIMEOUT); + + let output = cmd.output().expect("listing command should run"); + assert!( + output.status.success(), + "{binary} failed: {:?}\nstderr: {}", + output.status, + String::from_utf8_lossy(&output.stderr) + ); + assert!(output.stderr.is_empty(), "{binary} wrote to stderr"); + let stdout = String::from_utf8(output.stdout).expect("UTF-8 output"); + let names = stdout.lines().map(str::to_owned).collect::>(); + assert!(!names.is_empty(), "{binary} returned no variables"); + names +} + +fn list_config_env_vars(binary: &str) -> Vec { + let directory = tempfile::tempdir().expect("temporary directory"); + list_config_env_vars_in( + binary, + directory.path(), + std::iter::empty::<(&str, &str)>(), + &[], + ) +} + +#[test] +fn config_env_listing_exits_before_startup_for_each_binary() { + for (binary, config_env, dotenv_env, invalid_env) in [ + ( + "iggy-server", + "IGGY_CONFIG_PATH", + "IGGY_ENV_PATH", + "IGGY_HTTP_ADDRESS", + ), + ( + "iggy-connectors", + "IGGY_CONNECTORS_CONFIG_PATH", + "IGGY_CONNECTORS_ENV_PATH", + "IGGY_CONNECTORS_HTTP_ADDRESS", + ), + ( + "iggy-mcp", + "IGGY_MCP_CONFIG_PATH", + "IGGY_MCP_ENV_PATH", + "IGGY_MCP_HTTP_ADDRESS", + ), + ] { + let directory = tempfile::tempdir().expect("temporary directory"); + let config_path = directory.path().join("invalid-config.toml"); + let dotenv_path = directory.path().join("invalid.env"); + std::fs::write(&config_path, "not valid toml = [").expect("write invalid config"); + std::fs::write(&dotenv_path, "not a dotenv assignment").expect("write invalid dotenv"); + + let entries_before: HashSet<_> = std::fs::read_dir(directory.path()) + .expect("read temporary directory") + .map(|entry| entry.expect("directory entry").file_name()) + .collect(); + let names = list_config_env_vars_in( + binary, + directory.path(), + [ + (config_env, config_path.into_os_string()), + (dotenv_env, dotenv_path.into_os_string()), + (invalid_env, "invalid-value".into()), + ], + &[], + ); + + let baseline = list_config_env_vars(binary); + assert_eq!( + names, baseline, + "{binary} output changed with invalid config, dotenv or env values" + ); + let entries_after: HashSet<_> = std::fs::read_dir(directory.path()) + .expect("read temporary directory") + .map(|entry| entry.expect("directory entry").file_name()) + .collect(); + assert_eq!( + entries_after, entries_before, + "{binary} created startup side effects" + ); + + assert!( + names.windows(2).all(|pair| pair[0] < pair[1]), + "{binary} output must be sorted and deduplicated" + ); + assert!( + names.iter().all(|name| name.starts_with("IGGY_")), + "{binary} output contained non-IGGY_ prefixed names" + ); + } +} + +#[test] +fn config_env_listing_includes_vector_index_templates() { + let names = list_config_env_vars("iggy-server"); + let names: HashSet<_> = names.iter().map(String::as_str).collect(); + + // Cluster nodes should have indexed templates + assert!( + names + .iter() + .any(|name| name.contains("IGGY_CLUSTER_NODES__")), + "server should list cluster node index templates" + ); + + // Verify nested vector expansion (nested placeholders) + assert!( + names + .iter() + .any(|name| name.contains("IGGY_CLUSTER_NODES__ADVERTISED_ADDRESSES__")), + "server should list nested vector templates" + ); +} + +#[test] +fn config_env_listing_includes_connector_templates() { + let names = list_config_env_vars("iggy-connectors"); + let names: HashSet<_> = names.iter().map(String::as_str).collect(); + + assert!( + names.contains("IGGY_CONNECTORS_SINK__PLUGIN_CONFIG_"), + "connectors should list SINK plugin config templates" + ); + assert!( + names.contains("IGGY_CONNECTORS_SOURCE__PLUGIN_CONFIG_"), + "connectors should list SOURCE plugin config templates" + ); + assert!( + names.contains("IGGY_CONNECTORS_SINK__PLUGIN_CONFIG_FORMAT"), + "connectors should list the typed SINK plugin config format" + ); + assert!( + names.contains("IGGY_CONNECTORS_SOURCE__PLUGIN_CONFIG_FORMAT"), + "connectors should list the typed SOURCE plugin config format" + ); + + assert!( + names.contains("IGGY_CONNECTORS_SINK__ENABLED"), + "connectors should list SINK__ENABLED" + ); + assert!( + names.contains("IGGY_CONNECTORS_SOURCE__ENABLED"), + "connectors should list SOURCE__ENABLED" + ); + + assert_eq!( + names + .iter() + .filter(|name| **name == "IGGY_CONNECTORS_CONNECTORS_CONFIG_TYPE") + .count(), + 1, + "enum variants should produce one deduplicated tag name" + ); +} + +#[test] +fn config_env_listing_excludes_other_processes_variables() { + let server = list_config_env_vars("iggy-server"); + let server_names: HashSet<_> = server.iter().map(String::as_str).collect(); + for excluded in [ + "IGGY_CI_BUILD", + "IGGY_HOME", + "IGGY_PASSWORD", + "IGGY_TEST_CLEANUP_DISABLED", + "IGGY_TEST_CLUSTER_NODES", + "IGGY_TEST_VERBOSE", + "IGGY_USERNAME", + ] { + assert!(!server_names.contains(excluded), "server listed {excluded}"); + } + assert!( + server_names + .iter() + .all(|name| !name.starts_with("IGGY_CONNECTORS_") + && !name.starts_with("IGGY_KAFKA_") + && !name.starts_with("IGGY_MCP_")), + "server listed a sibling binary's variables" + ); + + for (binary, prefix) in [ + ("iggy-connectors", "IGGY_CONNECTORS_"), + ("iggy-mcp", "IGGY_MCP_"), + ] { + let names = list_config_env_vars(binary); + assert!( + names + .iter() + .all(|name| name == "IGGY_DISPLAY_CONFIG" || name.starts_with(prefix)), + "{binary} listed a variable outside {prefix}" + ); + } +} + +#[test] +fn config_env_listing_includes_each_process_runtime_variables() { + for (binary, expected) in [ + ( + "iggy-server", + &["IGGY_ROOT_PASSWORD", "IGGY_SHARD_RUNTIME_CAPACITY"][..], + ), + ( + "iggy-connectors", + &["IGGY_CONNECTORS_CONFIG_PATH", "IGGY_CONNECTORS_ENV_PATH"][..], + ), + ( + "iggy-mcp", + &["IGGY_MCP_CONFIG_PATH", "IGGY_MCP_ENV_PATH"][..], + ), + ] { + let names = list_config_env_vars(binary); + for expected_name in expected { + assert!( + names.iter().any(|name| name == expected_name), + "{binary} did not list {expected_name}" + ); + } + } +} + +#[test] +fn config_env_listing_with_fresh_does_not_wipe_data_dir() { + // Verify that --fresh does not wipe data when used with --list-config-env-vars + // Create a tempdir with a sentinel file + let directory = tempfile::tempdir().expect("temporary directory"); + let data_dir = directory.path().join("data"); + std::fs::create_dir(&data_dir).expect("create data directory"); + let sentinel = data_dir.join("sentinel.txt"); + std::fs::write(&sentinel, "sentinel content").expect("write sentinel file"); + + // Run with both --fresh and --list-config-env-vars + let output_with_fresh = list_config_env_vars_in( + "iggy-server", + directory.path(), + [("IGGY_PATH", data_dir.to_string_lossy().into_owned())], + &["--fresh"], + ); + + assert!( + sentinel.exists(), + "sentinel file should not be deleted by early exit" + ); + assert_eq!( + std::fs::read_to_string(&sentinel).expect("read sentinel"), + "sentinel content", + "sentinel file should not be modified" + ); + + let output_without_fresh = list_config_env_vars("iggy-server"); + assert_eq!( + output_with_fresh, output_without_fresh, + "output should be identical with and without --fresh" + ); +} diff --git a/core/integration/tests/mod.rs b/core/integration/tests/mod.rs index d34f2b00c8..ae40a2d4c7 100644 --- a/core/integration/tests/mod.rs +++ b/core/integration/tests/mod.rs @@ -37,6 +37,7 @@ mod cli; // Raw-wire spec tests for VSR session continuity across a node restart // (IGGY-137). mod cluster; +mod config_env_listing; mod config_provider; mod connectors; mod data_integrity; diff --git a/core/server/README.md b/core/server/README.md index 7c4d16147a..21348d419c 100644 --- a/core/server/README.md +++ b/core/server/README.md @@ -27,7 +27,7 @@ To run one node of a cluster, pass its replica ID from the `cluster.nodes` roste cargo run --bin iggy-server --release -- --replica-id 0 ``` -`--replica-id` is the only command line argument; everything else is configuration. +Command line arguments include `--replica-id`, `--fresh` (`-f`), `--with-default-root-credentials`, and `--list-config-env-vars`; everything else is configuration. ## Configuration diff --git a/core/server/src/args.rs b/core/server/src/args.rs index 9025c98f94..bd8a4c2341 100644 --- a/core/server/src/args.rs +++ b/core/server/src/args.rs @@ -87,6 +87,25 @@ For more information, visit: https://iggy.apache.org/docs/introduction/getting-s // variable names and paths must stay unquoted rather than wear rustdoc backticks. #[allow(clippy::doc_markdown)] pub struct Args { + #[arg( + long, + help = "Print supported configuration environment variables and exit", + long_help = r#"Print supported configuration environment variables and exit. + +Lists all supported IGGY_* environment variable names and templates, +sorted and deduplicated. Template syntax: +- represents vector indices (0-255 for most fields, except + cluster.nodes[*].advertised_addresses, which caps at 0-15) + +Exits immediately before any startup (before dotenv, config loading, +logging, runtimes, credentials, plugins, filesystem or network activity). +Works even with missing or invalid configuration files. + +Example: + iggy-server --list-config-env-vars"# + )] + pub list_config_env_vars: bool, + /// Remove the system path before starting (WARNING: THIS WILL DELETE ALL DATA!) /// /// Deletes the configured system data directory ('local_data' by default, diff --git a/core/server/src/main.rs b/core/server/src/main.rs index 75f8303b13..a0bc0473b8 100644 --- a/core/server/src/main.rs +++ b/core/server/src/main.rs @@ -22,7 +22,7 @@ mod banner; use args::Args; use clap::Parser; -use configs::server::ServerConfig; +use configs::{ConfigEnvMappings, print_env_var_names, server::ServerConfig}; use server::boot::{ apply_default_root_credentials, bootstrap, load_config, prepare_runtime_dirs, raise_open_file_limit, @@ -42,6 +42,10 @@ fn main() -> Result<(), ServerError> { // visible. `create_shard_executor` also reads its capacity knob from the // environment, which is why the `.env` load has to precede it. let args = Args::parse(); + if args.list_config_env_vars { + print_config_env_vars().map_err(ServerError::ListConfigEnvVars)?; + return Ok(()); + } banner::print(server::VERSION); // `logging` owns the tracing appender worker guards; it must outlive the // shard threads or every log line after bootstrap is silently dropped. @@ -132,3 +136,14 @@ fn main() -> Result<(), ServerError> { info!("server shutdown complete"); Ok(()) } + +fn print_config_env_vars() -> std::io::Result<()> { + let mut stdout = std::io::stdout(); + print_env_var_names( + ServerConfig::env_templates() + .iter() + .map(|t| t.env_name) + .chain(configs::server::server_runtime_env_vars()), + &mut stdout, + ) +} diff --git a/core/server/src/server_error.rs b/core/server/src/server_error.rs index 74c364c1e4..6ab9a6da90 100644 --- a/core/server/src/server_error.rs +++ b/core/server/src/server_error.rs @@ -339,6 +339,8 @@ pub enum ServerError { /// as clean. #[error("server shut down after a panic: {description}")] Panicked { description: String }, + #[error("Failed to list config environment variables")] + ListConfigEnvVars(#[source] std::io::Error), } /// Per-shard outcome captured by [`crate::boot::ShardHandles::join_all`]