diff --git a/flink-sql-runner/src/main/java/com/datasqrl/flinkrunner/BaseRunner.java b/flink-sql-runner/src/main/java/com/datasqrl/flinkrunner/BaseRunner.java index eeb9ccbd..f08553d4 100644 --- a/flink-sql-runner/src/main/java/com/datasqrl/flinkrunner/BaseRunner.java +++ b/flink-sql-runner/src/main/java/com/datasqrl/flinkrunner/BaseRunner.java @@ -91,7 +91,7 @@ Configuration initConfiguration() { var conf = new Configuration(); if (StringUtils.isNotBlank(configDir)) { log.info("Loading Flink configuration from '{}'", configDir); - conf = GlobalConfiguration.loadConfiguration(configDir); + conf = resolveEnvVars(GlobalConfiguration.loadConfiguration(configDir)); } // Do not overwrite runtime given in YAML @@ -102,6 +102,20 @@ Configuration initConfiguration() { return conf; } + private Configuration resolveEnvVars(Configuration source) { + var conf = new Configuration(source); + source + .toMap() + .forEach( + (key, value) -> { + var resolved = resolver.resolve(value); + if (!resolved.equals(value)) { + conf.setString(key, resolved); + } + }); + return conf; + } + static String readTextFile(String path) throws IOException { var f = new File(path); if (!f.exists()) { diff --git a/flink-sql-runner/src/test/java/com/datasqrl/flinkrunner/CliRunnerTest.java b/flink-sql-runner/src/test/java/com/datasqrl/flinkrunner/CliRunnerTest.java index d059acbc..9a65d4b6 100644 --- a/flink-sql-runner/src/test/java/com/datasqrl/flinkrunner/CliRunnerTest.java +++ b/flink-sql-runner/src/test/java/com/datasqrl/flinkrunner/CliRunnerTest.java @@ -26,6 +26,7 @@ import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; +import java.util.Map; import java.util.function.Supplier; import org.apache.flink.api.common.RuntimeExecutionMode; import org.apache.flink.configuration.ExecutionOptions; @@ -171,4 +172,27 @@ void initConfiguration_shouldCreateNewConfigIfNoPathGiven(RuntimeExecutionMode m assertThat(finalConf.keySet()).hasSize(1); assertThat(finalConf.get(ExecutionOptions.RUNTIME_MODE)).isEqualTo(mode); } + + @Test + void initConfiguration_shouldResolveEnvVarPlaceholdersInConfig() throws IOException { + // Arrange — a secret injected via env var, e.g. the Iceberg maintenance lock password + Path configFile = tempDir.resolve("config.yaml"); + Files.writeString( + configFile, + "flink-maintenance.lock.jdbc.password: ${ICEBERG_LOCK_PASSWORD}\ndummy.key: literal"); + + var resolver = + EnvVarResolver.builder().envVars(Map.of("ICEBERG_LOCK_PASSWORD", "s3cr3t")).build(); + var cliRunner = + new CliRunner( + RuntimeExecutionMode.STREAMING, resolver, null, null, tempDir.toString(), null); + + // Act + var finalConf = cliRunner.initConfiguration(); + + // Assert + assertThat(finalConf.toMap()) + .containsEntry("flink-maintenance.lock.jdbc.password", "s3cr3t") + .containsEntry("dummy.key", "literal"); + } }