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 221784105f9f4f..d03616998d6e1c 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 @@ protected CompletableFuture handleRequest( SessionHandle sessionHandle = request.getPathParameter(SessionHandleIdPathParameter.class); Map sessionConfig = this.service.getSessionConfig(sessionHandle); + Map 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 5735abd6dee318..f67b3d84130f82 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.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; @@ -623,7 +624,9 @@ private ResultFetcher callSetOperation( return ResultFetcher.fromTableResult(handle, TABLE_RESULT_OK, false); } else if (setOp.getKey().isEmpty() && setOp.getValue().isEmpty()) { // show all properties - Map configMap = tableEnv.getConfig().getConfiguration().toMap(); + Map 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 3892524669d100..70ba85d44023b1 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 @@ void testGetSessionConfiguration() throws Exception { } } + @Test + void testGetSessionConfigurationHidesSensitiveValues() throws Exception { + Map sensitiveProperties = new HashMap<>(); + sensitiveProperties.put("s3.secret-key", "super-secret-value"); + CompletableFuture openResponse = + sendRequest( + openSessionHeaders, + emptyParameters, + new OpenSessionRequestBody(SESSION_NAME, sensitiveProperties)); + SessionHandle handle = + new SessionHandle(UUID.fromString(openResponse.get().getSessionHandle())); + SessionMessageParameters parameters = new SessionMessageParameters(handle); + + CompletableFuture future = + sendRequest(GetSessionConfigHeaders.getInstance(), parameters, emptyRequestBody); + Map 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 a9dc23efd0ab3c..a9e9d57b7c0779 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 @@ 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 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 @@ void testResultKind() throws Exception { 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 rows = fetchAllResults(service, sessionHandle, setAllHandle); + + Map 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 void validateResultSetField( SessionHandle sessionHandle, String statement,