This is an automated email from the ASF dual-hosted git repository.
gaborgsomogyi pushed a commit to branch release-2.1
in repository https://gitbox.apache.org/repos/asf/flink.git
The following commit(s) were added to refs/heads/release-2.1 by this push:
new ff96c1d2dbe [FLINK-40539][sql-gateway] Redact sensitive values in SET;
and session config responses
ff96c1d2dbe is described below
commit ff96c1d2dbe4345d5a74185c60c8123416233775
Author: Gabor Somogyi <[email protected]>
AuthorDate: Fri Sep 4 07:57:26 2026 +0200
[FLINK-40539][sql-gateway] Redact sensitive values in SET; and session
config responses
---
.../handler/session/GetSessionConfigHandler.java | 5 +++-
.../service/operation/OperationExecutor.java | 5 +++-
.../table/gateway/rest/SessionRelatedITCase.java | 25 ++++++++++++++++++
.../service/SqlGatewayServiceStatementITCase.java | 30 ++++++++++++++++++++++
4 files changed, 63 insertions(+), 2 deletions(-)
diff --git
a/flink-table/flink-sql-gateway/src/main/java/org/apache/flink/table/gateway/rest/handler/session/GetSessionConfigHandler.java
b/flink-table/flink-sql-gateway/src/main/java/org/apache/flink/table/gateway/rest/handler/session/GetSessionConfigHandler.java
index 221784105f9..d03616998d6 100644
---
a/flink-table/flink-sql-gateway/src/main/java/org/apache/flink/table/gateway/rest/handler/session/GetSessionConfigHandler.java
+++
b/flink-table/flink-sql-gateway/src/main/java/org/apache/flink/table/gateway/rest/handler/session/GetSessionConfigHandler.java
@@ -18,6 +18,7 @@
package org.apache.flink.table.gateway.rest.handler.session;
+import org.apache.flink.configuration.ConfigurationUtils;
import org.apache.flink.runtime.rest.handler.HandlerRequest;
import org.apache.flink.runtime.rest.handler.RestHandlerException;
import org.apache.flink.runtime.rest.messages.EmptyRequestBody;
@@ -59,8 +60,10 @@ public class GetSessionConfigHandler
SessionHandle sessionHandle =
request.getPathParameter(SessionHandleIdPathParameter.class);
Map<String, String> sessionConfig =
this.service.getSessionConfig(sessionHandle);
+ Map<String, String> redactedSessionConfig =
+ ConfigurationUtils.hideSensitiveValues(sessionConfig);
return CompletableFuture.completedFuture(
- new GetSessionConfigResponseBody(sessionConfig));
+ new GetSessionConfigResponseBody(redactedSessionConfig));
} catch (SqlGatewayException e) {
throw new RestHandlerException(
e.getMessage(), HttpResponseStatus.INTERNAL_SERVER_ERROR,
e);
diff --git
a/flink-table/flink-sql-gateway/src/main/java/org/apache/flink/table/gateway/service/operation/OperationExecutor.java
b/flink-table/flink-sql-gateway/src/main/java/org/apache/flink/table/gateway/service/operation/OperationExecutor.java
index 9612b57fcca..aa4c5788fbe 100644
---
a/flink-table/flink-sql-gateway/src/main/java/org/apache/flink/table/gateway/service/operation/OperationExecutor.java
+++
b/flink-table/flink-sql-gateway/src/main/java/org/apache/flink/table/gateway/service/operation/OperationExecutor.java
@@ -28,6 +28,7 @@ import
org.apache.flink.client.deployment.DefaultClusterClientServiceLoader;
import org.apache.flink.client.program.ClusterClient;
import org.apache.flink.configuration.CheckpointingOptions;
import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.ConfigurationUtils;
import org.apache.flink.core.execution.SavepointFormatType;
import org.apache.flink.runtime.client.JobStatusMessage;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
@@ -610,7 +611,9 @@ public class OperationExecutor {
return ResultFetcher.fromTableResult(handle, TABLE_RESULT_OK,
false);
} else if (!setOp.getKey().isPresent() &&
!setOp.getValue().isPresent()) {
// show all properties
- Map<String, String> configMap =
tableEnv.getConfig().getConfiguration().toMap();
+ Map<String, String> configMap =
+ ConfigurationUtils.hideSensitiveValues(
+ tableEnv.getConfig().getConfiguration().toMap());
return ResultFetcher.fromResults(
handle,
ResolvedSchema.of(
diff --git
a/flink-table/flink-sql-gateway/src/test/java/org/apache/flink/table/gateway/rest/SessionRelatedITCase.java
b/flink-table/flink-sql-gateway/src/test/java/org/apache/flink/table/gateway/rest/SessionRelatedITCase.java
index 3892524669d..70ba85d4402 100644
---
a/flink-table/flink-sql-gateway/src/test/java/org/apache/flink/table/gateway/rest/SessionRelatedITCase.java
+++
b/flink-table/flink-sql-gateway/src/test/java/org/apache/flink/table/gateway/rest/SessionRelatedITCase.java
@@ -18,6 +18,7 @@
package org.apache.flink.table.gateway.rest;
+import org.apache.flink.configuration.GlobalConfiguration;
import org.apache.flink.runtime.rest.messages.EmptyMessageParameters;
import org.apache.flink.runtime.rest.messages.EmptyRequestBody;
import org.apache.flink.runtime.rest.messages.EmptyResponseBody;
@@ -158,6 +159,30 @@ class SessionRelatedITCase extends RestAPIITCaseBase {
}
}
+ @Test
+ void testGetSessionConfigurationHidesSensitiveValues() throws Exception {
+ Map<String, String> sensitiveProperties = new HashMap<>();
+ sensitiveProperties.put("s3.secret-key", "super-secret-value");
+ CompletableFuture<OpenSessionResponseBody> openResponse =
+ sendRequest(
+ openSessionHeaders,
+ emptyParameters,
+ new OpenSessionRequestBody(SESSION_NAME,
sensitiveProperties));
+ SessionHandle handle =
+ new
SessionHandle(UUID.fromString(openResponse.get().getSessionHandle()));
+ SessionMessageParameters parameters = new
SessionMessageParameters(handle);
+
+ CompletableFuture<GetSessionConfigResponseBody> future =
+ sendRequest(GetSessionConfigHeaders.getInstance(), parameters,
emptyRequestBody);
+ Map<String, String> getProperties = future.get().getProperties();
+
+ assertThat(getProperties).containsKey("s3.secret-key");
+ assertThat(getProperties.get("s3.secret-key"))
+ .isEqualTo(GlobalConfiguration.HIDDEN_CONTENT);
+
+ sendRequest(closeSessionHeaders, parameters, emptyRequestBody).get();
+ }
+
@Test
void testTouchSession() throws Exception {
Session session =
diff --git
a/flink-table/flink-sql-gateway/src/test/java/org/apache/flink/table/gateway/service/SqlGatewayServiceStatementITCase.java
b/flink-table/flink-sql-gateway/src/test/java/org/apache/flink/table/gateway/service/SqlGatewayServiceStatementITCase.java
index a9dc23efd0a..a9e9d57b7c0 100644
---
a/flink-table/flink-sql-gateway/src/test/java/org/apache/flink/table/gateway/service/SqlGatewayServiceStatementITCase.java
+++
b/flink-table/flink-sql-gateway/src/test/java/org/apache/flink/table/gateway/service/SqlGatewayServiceStatementITCase.java
@@ -21,6 +21,7 @@ package org.apache.flink.table.gateway.service;
import org.apache.flink.api.common.RuntimeExecutionMode;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.ExecutionOptions;
+import org.apache.flink.configuration.GlobalConfiguration;
import org.apache.flink.core.testutils.CommonTestUtils;
import org.apache.flink.table.api.ResultKind;
import org.apache.flink.table.data.RowData;
@@ -46,14 +47,17 @@ import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Path;
import java.time.Duration;
import java.util.Collections;
+import java.util.HashMap;
import java.util.Iterator;
import java.util.List;
+import java.util.Map;
import java.util.function.BiFunction;
import java.util.stream.Collectors;
import static
org.apache.flink.table.gateway.api.config.SqlGatewayServiceConfigOptions.SQL_GATEWAY_SESSION_PLAN_CACHE_ENABLED;
import static
org.apache.flink.table.gateway.service.utils.SqlGatewayServiceTestUtil.awaitOperationTermination;
import static
org.apache.flink.table.gateway.service.utils.SqlGatewayServiceTestUtil.createInitializedSession;
+import static
org.apache.flink.table.gateway.service.utils.SqlGatewayServiceTestUtil.fetchAllResults;
import static
org.apache.flink.table.gateway.service.utils.SqlGatewayServiceTestUtil.fetchResults;
import static org.assertj.core.api.Assertions.assertThat;
@@ -286,6 +290,32 @@ public class SqlGatewayServiceStatementITCase extends
AbstractSqlGatewayStatemen
sessionHandle, "SET;", resultKindGetter,
ResultKind.SUCCESS_WITH_CONTENT);
}
+ @Test
+ void testSetHidesSensitiveValues() throws Exception {
+ SessionHandle sessionHandle = createInitializedSession(service);
+
+ runAndAwait(sessionHandle, "SET 's3.secret-key' =
'super-secret-value';");
+
+ OperationHandle setAllHandle = runAndAwait(sessionHandle, "SET;");
+ List<RowData> rows = fetchAllResults(service, sessionHandle,
setAllHandle);
+
+ Map<String, String> properties = new HashMap<>();
+ for (RowData row : rows) {
+ properties.put(row.getString(0).toString(),
row.getString(1).toString());
+ }
+
+ assertThat(properties).containsKey("s3.secret-key");
+
assertThat(properties.get("s3.secret-key")).isEqualTo(GlobalConfiguration.HIDDEN_CONTENT);
+ }
+
+ private OperationHandle runAndAwait(SessionHandle sessionHandle, String
statement)
+ throws Exception {
+ OperationHandle operationHandle =
+ service.executeStatement(sessionHandle, statement, -1, new
Configuration());
+ awaitOperationTermination(service, sessionHandle, operationHandle);
+ return operationHandle;
+ }
+
private <T> void validateResultSetField(
SessionHandle sessionHandle,
String statement,